本文目录导读:

我来为您详细介绍如何制定Python数据迭代同步规范。
基础迭代器规范
from typing import Iterator, Iterable, TypeVar, Generic
from abc import ABC, abstractmethod
import logging
T = TypeVar('T')
class DataIterable(Generic[T], ABC):
"""数据迭代器规范基类"""
@abstractmethod
def __iter__(self) -> Iterator[T]:
"""返回迭代器"""
pass
@abstractmethod
def __len__(self) -> int:
"""返回数据长度"""
pass
@abstractmethod
def validate(self) -> bool:
"""数据验证方法"""
pass
同步规范定义
from dataclasses import dataclass, field
from datetime import datetime, timedelta
from enum import Enum
import hashlib
import json
class SyncStatus(Enum):
"""同步状态枚举"""
PENDING = "pending"
SYNCING = "syncing"
COMPLETED = "completed"
FAILED = "failed"
PARTIAL = "partial"
@dataclass
class SyncRecord:
"""同步记录"""
source: str
target: str
batch_id: str
status: SyncStatus
records_count: int
start_time: datetime
end_time: datetime = None
error_message: str = None
checksum: str = None
def generate_checksum(self, data: list) -> str:
"""生成数据校验和"""
data_str = json.dumps(data, sort_keys=True)
return hashlib.md5(data_str.encode()).hexdigest()
class SyncSpecification:
"""同步规范定义"""
def __init__(self,
batch_size: int = 1000,
retry_times: int = 3,
timeout: int = 30,
max_retry_delay: int = 60):
self.batch_size = batch_size
self.retry_times = retry_times
self.timeout = timeout
self.max_retry_delay = max_retry_delay
self.sync_records = []
def validate_sync_params(self):
"""验证同步参数"""
assert self.batch_size > 0, "批次大小必须大于0"
assert self.retry_times >= 0, "重试次数不能为负数"
assert self.timeout > 0, "超时时间必须大于0"
数据迭代器实现
from typing import List, Optional, Generator
import time
from contextlib import contextmanager
class BatchIterator:
"""分批迭代器实现"""
def __init__(self, data: List, batch_size: int = 1000):
self.data = data
self.batch_size = batch_size
self.index = 0
def __iter__(self):
return self
def __next__(self):
if self.index >= len(self.data):
raise StopIteration
batch = self.data[self.index:self.index + self.batch_size]
self.index += self.batch_size
return batch
def __len__(self):
return (len(self.data) + self.batch_size - 1) // self.batch_size
class SafeIterator:
"""安全迭代器(带错误处理)"""
def __init__(self, data_source, max_retries: int = 3):
self.data_source = data_source
self.max_retries = max_retries
self.current_retry = 0
def __iter__(self):
return self
def __next__(self):
while self.current_retry < self.max_retries:
try:
data = next(self.data_source)
self.current_retry = 0
return data
except StopIteration:
raise
except Exception as e:
self.current_retry += 1
if self.current_retry >= self.max_retries:
raise RuntimeError(f"迭代失败,已重试{self.max_retries}次") from e
time.sleep(2 ** self.current_retry) # 指数退避
同步管理器
from concurrent.futures import ThreadPoolExecutor, as_completed
import threading
from typing import Callable, Any
class SyncManager:
"""同步管理器"""
def __init__(self, spec: SyncSpecification):
self.spec = spec
self.logger = logging.getLogger(__name__)
self.lock = threading.Lock()
self.progress = 0
def sync_data(self,
source_iterator: Iterator,
target_writer: Callable,
progress_callback: Callable = None) -> SyncRecord:
"""执行数据同步"""
batch_id = f"sync_{datetime.now().strftime('%Y%m%d_%H%M%S')}"
record = SyncRecord(
source=str(source_iterator),
target=str(target_writer),
batch_id=batch_id,
status=SyncStatus.SYNCING,
records_count=0,
start_time=datetime.now()
)
try:
total_processed = 0
for batch in source_iterator:
# 处理批次数据
processed_data = self._process_batch(batch)
# 写入目标
target_writer(processed_data)
# 更新进度
total_processed += len(batch)
self.progress = total_processed
if progress_callback:
progress_callback(total_processed)
# 速率限制
time.sleep(0.1)
# 更新同步记录
record.status = SyncStatus.COMPLETED
record.records_count = total_processed
record.end_time = datetime.now()
record.checksum = record.generate_checksum(processed_data)
except Exception as e:
record.status = SyncStatus.FAILED
record.error_message = str(e)
self.logger.error(f"同步失败: {e}")
finally:
self.spec.sync_records.append(record)
return record
def _process_batch(self, batch: List) -> List:
"""处理批次数据(可重写)"""
return batch
def get_progress(self) -> float:
"""获取同步进度"""
with self.lock:
return self.progress
并发同步控制
import asyncio
from typing import AsyncIterator
class AsyncSyncManager:
"""异步同步管理器"""
def __init__(self, max_concurrent: int = 5):
self.max_concurrent = max_concurrent
self.semaphore = asyncio.Semaphore(max_concurrent)
async def concurrent_sync(self,
sources: List[AsyncIterator],
target: Callable) -> List[SyncRecord]:
"""并发同步多个数据源"""
async def sync_single(source: AsyncIterator) -> SyncRecord:
async with self.semaphore:
return await self._sync_source(source, target)
tasks = [sync_single(source) for source in sources]
return await asyncio.gather(*tasks, return_exceptions=True)
async def _sync_source(self,
source: AsyncIterator,
target: Callable) -> SyncRecord:
"""同步单个数据源"""
record = SyncRecord(
source=str(source),
target=str(target),
batch_id=f"async_{datetime.now().timestamp()}",
status=SyncStatus.SYNCING,
records_count=0,
start_time=datetime.now()
)
try:
async for batch in source:
await target(batch)
record.records_count += len(batch)
record.status = SyncStatus.COMPLETED
record.end_time = datetime.now()
except Exception as e:
record.status = SyncStatus.FAILED
record.error_message = str(e)
return record
数据验证规范
from pydantic import BaseModel, validator
from typing import Optional
class SyncDataModel(BaseModel):
"""同步数据模型"""
id: str
timestamp: datetime
data: dict
valid: bool = True
@validator('id')
def validate_id(cls, v):
assert len(v) > 0, "ID不能为空"
return v
class DataValidator:
"""数据验证器"""
def __init__(self):
self.validation_rules = []
def add_rule(self, rule: Callable):
"""添加验证规则"""
self.validation_rules.append(rule)
def validate(self, data: List[Dict]) -> tuple[List, List]:
"""验证数据,返回有效和无效数据"""
valid_data = []
invalid_data = []
for item in data:
try:
validated = SyncDataModel(**item)
valid_data.append(validated.dict())
except Exception as e:
item['error'] = str(e)
invalid_data.append(item)
return valid_data, invalid_data
使用示例
def main():
# 创建同步规范
spec = SyncSpecification(
batch_size=500,
retry_times=3,
timeout=60
)
# 创建同步管理器
manager = SyncManager(spec)
# 准备测试数据
test_data = [{"id": i, "name": f"item_{i}"} for i in range(10000)]
# 创建迭代器
batch_iterator = BatchIterator(test_data, batch_size=500)
# 定义写入函数
def write_to_target(data):
# 模拟写入目标
print(f"写入{len(data)}条数据")
# 实际应用中这里连接数据库或API
# 执行同步
result = manager.sync_data(
source_iterator=batch_iterator,
target_writer=write_to_target,
progress_callback=lambda x: print(f"进度: {x}")
)
print(f"同步结果: {result.status}")
print(f"处理记录数: {result.records_count}")
if __name__ == "__main__":
main()
最佳实践建议
规范要点:
- 错误处理:实现重试机制和降级策略
- 性能监控:记录同步时间和处理速率
- 数据完整性:使用校验和验证
- 日志记录:详细记录同步过程
- 幂等性:确保重复执行结果一致
- 失败恢复:支持断点续传
使用建议:
- 根据数据量和网络状况调整batch_size
- 使用异步IO提高大量数据的同步效率
- 实现监控告警机制
- 添加数据一致性检查
- 定期进行数据校验和修复
这套规范可以根据具体业务需求进行调整和扩展。