解决DynamoDB数据问题:生产环境JSON解析异常求助
批量修复DynamoDB中JSON格式异常记录方案
针对生产环境中DynamoDB百万级记录的JSON解析问题,以下是可落地的批量处理方案,无需人工逐条修复:
1. 批量识别异常记录
先通过扫描排查所有存在JSON解析问题的记录,将其主键存入临时表以便后续处理。推荐用Python+Boto3实现,利用DynamoDB的分页扫描避免超时,同时支持并行扫描提升效率(百万级数据建议开启)。
示例代码:
import boto3 import json dynamodb = boto3.resource('dynamodb') source_table = dynamodb.Table('your-production-table') invalid_table = dynamodb.Table('invalid-records-table') def scan_and_mark_invalid(): last_evaluated_key = None while True: scan_kwargs = {'Limit': 100} # 按需调整批次大小 if last_evaluated_key: scan_kwargs['ExclusiveStartKey'] = last_evaluated_key response = source_table.scan(**scan_kwargs) items = response.get('Items', []) invalid_items = [] for item in items: try: # 替换为你需要解析的JSON字段 json.loads(item['transaction_data']) except json.JSONDecodeError: # 仅存储主键和错误类型,减少临时表存储压力 invalid_items.append({ 'PutRequest': { 'Item': { 'id': item['id'], # 替换为你的表主键 'error_type': 'JSONDecodeError' } } }) if invalid_items: invalid_table.batch_write_item(RequestItems={'invalid-records-table': invalid_items}) last_evaluated_key = response.get('LastEvaluatedKey') if not last_evaluated_key: break if __name__ == '__main__': scan_and_mark_invalid()
2. 批量修复异常记录
分析JSON解析失败的具体原因(比如引号未闭合、转义错误、字段缺失等),编写针对性修复逻辑,从临时表读取异常记录主键,批量查询原表数据并修复后写回。
示例代码(针对未闭合引号的修复):
import re import boto3 import json dynamodb = boto3.resource('dynamodb') source_table = dynamodb.Table('your-production-table') invalid_table = dynamodb.Table('invalid-records-table') def fix_json_string(s): # 修复未转义的引号 s = re.sub(r'(?<!\\)"(?!:|,|\]|\}|$)', '\\"', s) # 补全未闭合的引号 if s.count('"') % 2 != 0: s += '"' return s def repair_invalid_records(): last_evaluated_key = None while True: scan_kwargs = {'Limit': 25} # 匹配BatchWriteItem的25条限制 if last_evaluated_key: scan_kwargs['ExclusiveStartKey'] = last_evaluated_key response = invalid_table.scan(**scan_kwargs) invalid_items = response.get('Items', []) if not invalid_items: break # 批量查询原表的异常记录 keys = [{'id': item['id']} for item in invalid_items] get_response = source_table.batch_get_item(RequestItems={'your-production-table': {'Keys': keys}}) items_to_repair = get_response['Responses']['your-production-table'] repair_requests = [] for item in items_to_repair: try: fixed_data = fix_json_string(item['transaction_data']) json.loads(fixed_data) # 验证修复后的合法性 item['transaction_data'] = fixed_data repair_requests.append({ 'PutRequest': { 'Item': item } }) except Exception as e: # 记录无法自动修复的记录,后续仅需处理少量此类数据 print(f"无法修复记录 {item['id']}: {str(e)}") if repair_requests: source_table.batch_write_item(RequestItems={'your-production-table': repair_requests}) last_evaluated_key = response.get('LastEvaluatedKey') if not last_evaluated_key: break if __name__ == '__main__': repair_invalid_records()
3. 修复夜间导入流程,杜绝后续问题
在夜间导入环节加入前置校验,确保只有合法JSON数据写入DynamoDB:
- 对数据源做JSON格式校验,过滤或修复异常数据后再入库;
- 加入异常监控,当无效记录数超过阈值时触发告警(比如集成SNS);
- 可选:用JSON Schema做结构化校验,确保数据字段类型、结构符合预期。
示例导入校验逻辑:
def validate_and_import(data): valid_items = [] invalid_count = 0 for item in data: try: json.loads(item['transaction_data']) valid_items.append(item) except json.JSONDecodeError: invalid_count += 1 print(f"无效JSON记录: {item['id']}") # 批量写入有效数据 if valid_items: source_table.batch_write_item(RequestItems={'your-production-table': [{'PutRequest': {'Item': item}} for item in valid_items]}) # 无效记录过多时触发告警 if invalid_count > 100: # 这里替换为你的告警逻辑,比如SNS发送通知 print(f"告警: 本次导入发现 {invalid_count} 条无效记录")
4. 优化数据提取流程的容错机制
在提取流程中加入异常处理,遇到JSON解析失败的记录时跳过并记录日志,不影响整体任务执行:
def extract_records(): last_evaluated_key = None while True: scan_kwargs = {'Limit': 100} if last_evaluated_key: scan_kwargs['ExclusiveStartKey'] = last_evaluated_key response = source_table.scan(**scan_kwargs) items = response.get('Items', []) for item in items: try: transaction_data = json.loads(item['transaction_data']) # 处理正常数据的逻辑 process_data(transaction_data) except json.JSONDecodeError: # 记录错误日志,后续可统一处理 with open('extraction_errors.log', 'a') as f: f.write(f"解析失败记录 {item['id']}: {item['transaction_data']}\n") last_evaluated_key = response.get('LastEvaluatedKey') if not last_evaluated_key: break
内容的提问来源于stack exchange,提问作者user1657861
相关产品推荐
相关产品推荐

