You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

解决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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.25 17:33:25