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})
排查结论及修复建议
核心问题:代码逻辑缺陷
- 变量名拼写错误:
update_tables函数中遍历buffers[topic_name],但参数传入的是top_name,直接导致批量写入时找不到对应缓冲区,这是旧表无法写入的关键原因。 - 缩进错误触发过早退出:
save_to_db中if counter>=30: exit()的缩进位置错误,它处于try块外,会在处理第30条数据前直接终止程序,后续数据完全无法处理。 - 未定义变量导致隐性异常:代码中直接使用
current_date_time、current_date、current_timestamp但未初始化,会触发未捕获的NameError,导致写入流程中断。 - 异常捕获无有效信息:
update_tables中的异常打印语句未填充变量,即使写入失败也无法定位具体原因。 - 去重逻辑不严谨:全局变量
received_values仅在进程内存中存储,程序重启后去重规则失效,且长期运行会导致内存占用过高。
AWS表配置排查方向
- 主键冲突:检查DDB表的主键设计,若写入数据的主键与存量数据重复,
put_item会直接覆盖旧数据,表现为“未写入”。 - 热点键限流:即使调整了WCU/RCU,若某几个
tob_id的写入过于集中,会触发隐性限流,按需模式下也可能出现此类情况,需查看CloudWatch的WriteThrottleEvents指标。 - 权限验证:确认程序使用的IAM角色是否拥有DDB表的
PutItem权限,旧表可能存在权限配置差异。
修复步骤
- 修正
update_tables中的变量名:将buffers[topic_name]改为buffers[top_name]。 - 调整退出逻辑的缩进:删除
save_to_db中重复的退出代码,保留update_tables中的退出逻辑。 - 初始化时间变量:在追加缓冲区前添加代码(需导入
datetime模块):current_date_time = datetime.datetime.now().isoformat() current_date = datetime.date.today().isoformat() current_timestamp = datetime.datetime.now().strftime('%H:%M:%S') - 补全异常打印信息:
print(f"Couldn't load data into table {table_bms.name}. Here's why: {err.response['Error']['Code']}: {err.response['Error']['Message']}") - 避免主键覆盖:若需新增而非覆盖数据,在
put_item中添加条件表达式:writer.put_item(Item=item, ConditionExpression='attribute_not_exists(tob_id) AND attribute_not_exists(timestamp)')
内容的提问来源于stack exchange,提问作者Sam
相关产品推荐
相关产品推荐

