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

DynamoDB存量表写入异常排查:3-4万条后无法正常插入

问题描述

代码功能:从Firebase获取数据,经去重、拆分后追加至缓冲区;当缓冲区达到4条限制时,使用batchWriter批量写入DynamoDB(DDB)表。

测试异常:

  • 新表可成功插入32条数据;
  • 表中存量达3-4万条后,无法插入所有预期数据;
  • 切换至已有60485条数据的旧表时,32条数据无一条写入,长时间执行仅插入1条;
  • 调整WCU/RCU为25或改用按需模式后问题仍存在,无报错及CloudWatch告警。

需排查问题源于代码逻辑还是AWS表配置。


代码1

# Define the DynamoDB tables
dynamodb = boto3.resource('dynamodb', region_name='my-region')
table_bms = dynamodb.Table('MS_TR_MS')

# Buffer and counter initialization
buffers = {
    'ms_str_msg': [],

}
counter = 0
buffer_size = 4
received_values = {}


def save_to_db(tob_id, top_name, top_data):
    with open(log_file_path, 'a') as log_file:
        try:
            # Check for duplicate entry
            key = f'{tob_id}_{top_name}'
            value = top_data
            if key in received_values and received_values[key] == value:
                # print(f"Duplicate entry for {key}: {value}. Discarding.")
                return
            else:
                received_values[key] = value
                print(f"Adding entry: {key}: {value}")
                
            # Split & Append data to the buffer            
            if isinstance(top_data, str):
                topic_data = top_data.split(":")[0]
                string_array = top_data.split(",")
            else:                
                return
                
        except Exception as e:
            print(f'Error: {top_data} \n{e}')
            return
# checking number of items inserted fro buffer to DDB table
        if counter>=30:
            exit()

    #Updating items to single table
        if top_name == "ms_str_msg":
            if len(string_array) == 5:
                # table = table_ms.name
                print('Updating MS table')

            buffers[top_name].append({
                # 'UID': {'S': unique_constant_value},
                "timestamp": str(current_date_time),
                "date": str(current_date),
                "time": str(current_timestamp),
                "tob_id": str(tob_id),
                "voltage": str(string_array[0]),
                "current_amps": str(string_array[1]),
                "total_capacity": str(string_array[2]),
                "SOC": str(string_array[3]),
                "Battery Level": str(string_array[4])})

        # Check buffer size and update tables
        if len(buffers[top_name]) >= buffer_size:
            update_tables(top_name)

# using batch_writer to write items from buffer to table

def update_tables(top_name):
    global counter
    try:
        with table_bms.batch_writer() as writer:
            for item in buffers[topic_name]:
                writer.put_item(Item=item)
                counter = counter+1
        buffers[top_name] = []
        
        # Adding a counter to check number of rows inserting
        print(f"Counter is {counter}")
        if counter>=30:
            exit()
            
    except ClientError as err:
        print("Couldn't load data into table %s. Here's why: %s: %s")

if __name__ == "__main__":
    main()

代码2

def update_tables():
    try:
        eastern = dateutil.tz.gettz('US/Eastern')
        with table_bms.batch_writer() as writer:
            for i in range(30000):
                msg = f"counter {i}"
                writer.put_item(Item={"Time": datetime.datetime.now(tz=eastern).strftime('%m/%d/%Y %H:%M:%S.%f %Z'), "Message": msg})

排查结论及修复建议

核心问题:代码逻辑缺陷

  1. 变量名拼写错误:update_tables函数中遍历buffers[topic_name],但参数传入的是top_name,直接导致批量写入时找不到对应缓冲区,这是旧表无法写入的关键原因。
  2. 缩进错误触发过早退出:save_to_db中if counter>=30: exit()的缩进位置错误,它处于try块外,会在处理第30条数据前直接终止程序,后续数据完全无法处理。
  3. 未定义变量导致隐性异常:代码中直接使用current_date_time、current_date、current_timestamp但未初始化,会触发未捕获的NameError,导致写入流程中断。
  4. 异常捕获无有效信息:update_tables中的异常打印语句未填充变量,即使写入失败也无法定位具体原因。
  5. 去重逻辑不严谨:全局变量received_values仅在进程内存中存储,程序重启后去重规则失效,且长期运行会导致内存占用过高。

AWS表配置排查方向

  1. 主键冲突:检查DDB表的主键设计,若写入数据的主键与存量数据重复,put_item会直接覆盖旧数据,表现为“未写入”。
  2. 热点键限流:即使调整了WCU/RCU,若某几个tob_id的写入过于集中,会触发隐性限流,按需模式下也可能出现此类情况,需查看CloudWatch的WriteThrottleEvents指标。
  3. 权限验证:确认程序使用的IAM角色是否拥有DDB表的PutItem权限,旧表可能存在权限配置差异。

修复步骤

  1. 修正update_tables中的变量名:将buffers[topic_name]改为buffers[top_name]。
  2. 调整退出逻辑的缩进:删除save_to_db中重复的退出代码,保留update_tables中的退出逻辑。
  3. 初始化时间变量:在追加缓冲区前添加代码(需导入datetime模块):
    current_date_time = datetime.datetime.now().isoformat()
    current_date = datetime.date.today().isoformat()
    current_timestamp = datetime.datetime.now().strftime('%H:%M:%S')
    
  4. 补全异常打印信息:
    print(f"Couldn't load data into table {table_bms.name}. Here's why: {err.response['Error']['Code']}: {err.response['Error']['Message']}")
    
  5. 避免主键覆盖:若需新增而非覆盖数据,在put_item中添加条件表达式:
    writer.put_item(Item=item, ConditionExpression='attribute_not_exists(tob_id) AND attribute_not_exists(timestamp)')
    

内容的提问来源于stack exchange,提问作者Sam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 19:34:57