1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
| # 实时 ETL 管道
from datetime import datetime, timedelta
import json
class StreamingETL:
def __init__(self, config):
self.config = config
self.validators = []
self.transformers = []
self.enrichers = []
def add_validator(self, validator):
"""添加数据验证器"""
self.validators.append(validator)
return self
def add_transformer(self, transformer):
"""添加数据转换器"""
self.transformers.append(transformer)
return self
def add_enricher(self, enricher):
"""添加数据增强器"""
self.enrichers.append(enricher)
return self
async def process_stream(self, stream):
"""处理数据流"""
async for batch in stream:
# 验证
valid_batch = await self.validate(batch)
# 转换
transformed_batch = await self.transform(valid_batch)
# 增强
enriched_batch = await self.enrich(transformed_batch)
# 输出
yield enriched_batch
async def validate(self, batch):
"""验证数据批次"""
valid_records = []
for record in batch:
is_valid = True
errors = []
for validator in self.validators:
result = await validator.validate(record)
if not result.is_valid:
is_valid = False
errors.extend(result.errors)
if is_valid:
valid_records.append(record)
else:
# 记录无效数据
await self.handle_invalid_record(record, errors)
return valid_records
async def transform(self, batch):
"""转换数据批次"""
transformed = batch
for transformer in self.transformers:
transformed = await transformer.transform(transformed)
return transformed
async def enrich(self, batch):
"""增强数据批次"""
enriched = []
for record in batch:
enriched_record = record.copy()
for enricher in self.enrichers:
enriched_record = await enricher.enrich(enriched_record)
enriched.append(enriched_record)
return enriched
async def handle_invalid_record(self, record, errors):
"""处理无效记录"""
error_record = {
'original_record': record,
'errors': errors,
'timestamp': datetime.utcnow().isoformat(),
'status': 'validation_failed'
}
# 发送到错误主题
await self.send_to_dlq(error_record)
async def send_to_dlq(self, record):
"""发送到死信队列"""
# 实现死信队列逻辑
pass
# 数据验证器
class SchemaValidator:
def __init__(self, schema):
self.schema = schema
async def validate(self, record):
"""验证记录是否符合 schema"""
errors = []
for field, rules in self.schema.items():
if field not in record:
if rules.get('required', False):
errors.append(f"Missing required field: {field}")
else:
value = record[field]
# 类型检查
if 'type' in rules:
if not isinstance(value, rules['type']):
errors.append(f"Field {field} has wrong type")
# 范围检查
if 'min' in rules and value < rules['min']:
errors.append(f"Field {field} below minimum")
if 'max' in rules and value > rules['max']:
errors.append(f"Field {field} above maximum")
return ValidationResult(
is_valid=len(errors) == 0,
errors=errors
)
# 数据转换器
class DataTransformer:
async def transform(self, batch):
"""转换数据批次"""
transformed = []
for record in batch:
transformed_record = {
'user_id': self.normalize_user_id(record.get('user_id')),
'timestamp': self.parse_timestamp(record.get('timestamp')),
'value': self.convert_value(record.get('value')),
'metadata': self.extract_metadata(record)
}
transformed.append(transformed_record)
return transformed
def normalize_user_id(self, user_id):
"""标准化用户 ID"""
if isinstance(user_id, str):
return user_id.strip().lower()
return str(user_id)
def parse_timestamp(self, timestamp):
"""解析时间戳"""
if isinstance(timestamp, (int, float)):
return datetime.fromtimestamp(timestamp)
elif isinstance(timestamp, str):
return datetime.fromisoformat(timestamp.replace('Z', '+00:00'))
return datetime.utcnow()
def convert_value(self, value):
"""转换数值"""
try:
return float(value)
except (TypeError, ValueError):
return 0.0
def extract_metadata(self, record):
"""提取元数据"""
return {
'source': record.get('source', 'unknown'),
'version': record.get('version', '1.0'),
'processed_at': datetime.utcnow().isoformat()
}
class ValidationResult:
def __init__(self, is_valid, errors):
self.is_valid = is_valid
self.errors = errors
|