AWS Glue写入DynamoDB遇重复键错误,如何定位重复项?
解决Glue导入DynamoDB时的重复键排查问题
一、直接识别重复的复合键数据
DynamoDB的重复键错误源于Partition Key + Sort Key的组合重复,你需要在写入DynamoDB前就排查出这些重复项——之前你打印的是写入后的结果框架,里面可能已经过滤了错误数据,所以看不到问题。
具体实现代码
假设你的Partition Key字段是pk,Sort Key字段是sk,在ApplyMapping之后、写入DynamoDB之前添加以下代码:
# 将动态框架转换为Spark DataFrame df = ApplyMapping_acuerdos_pago.toDF() # 按复合键分组,统计重复次数 duplicate_records = df.groupBy("pk", "sk").count().filter("count > 1") # 打印重复键(数据量大的话用limit限制输出条数,避免日志溢出) logger.info("检测到的重复复合键:") duplicate_records.show(truncate=False, limit=50) # 可选:将重复数据写入S3存储桶,方便后续分析 duplicate_records.write.mode("overwrite").csv("s3://your-bucket-name/duplicate-records/acuerdos_pago/")
二、捕获Glue发起的DynamoDB请求
Glue默认不会打印底层API请求内容,但可以通过两种方式实现:
1. 开启Glue作业DEBUG级日志
在Glue作业的配置页面,把日志级别调整为DEBUG,这样会打印包括DynamoDB批量写入请求在内的底层细节。注意:DEBUG日志会产生大量数据,排查完成后记得改回INFO级别。
2. 自定义批量写入逻辑,手动打印批次数据
放弃Glue内置的write_dynamic_frame.from_options,改用Spark的DynamoDB连接器手动实现写入,这样可以在每个批次发送前打印数据:
# 先将动态框架转为DataFrame df = ApplyMapping_acuerdos_pago.toDF() # 定义批量写入函数,加入日志打印 def process_batch(batch_df, batch_id): logger.info(f"正在处理批次 {batch_id},数据内容:") batch_df.show(truncate=False, limit=20) # 写入DynamoDB batch_df.write \ .format("dynamodb") \ .option("tableName", f"powerbi_sinco_{compamy_id}_acuerdos_pago_{environment}") \ .mode("append") \ .save() # 按批次处理数据(可根据需求调整批次大小) df.foreachBatch(process_batch)
三、额外排查技巧
- 直接在SQLServer源端排查:用SQL语句提前找出重复的复合键
SELECT [pk_column], [sk_column], COUNT(*) AS duplicate_count FROM adi_dtm.acuerdos_pago GROUP BY [pk_column], [sk_column] HAVING COUNT(*) > 1;
- 验证映射结果:检查
ApplyMapping后的字段是否正确,避免因字段类型转换、映射错误导致的“伪重复”
内容的提问来源于stack exchange,提问作者Jefferson
相关产品推荐
相关产品推荐

