Flink Sink写入重复项问题:如何排查DynamoDB写入异常记录?
解决DynamoDB Sink写入重复主键异常的排查与自定义处理方案
一、异常发生后检查写入记录的方法
- 开启详细日志追踪:在Sink的配置或代码中开启DEBUG级别的日志,将待写入记录的主键字段(如
partitionKey+sortKey组合)在写入前完整打印。当异常触发时,回溯日志即可定位对应批次里的重复主键记录。 - 暂存待写入批次数据:在写入DynamoDB之前,将整批次数据暂存到临时存储(如本地文件、内存队列),写入成功则清理,写入失败则保留这批数据,直接分析其中的主键重复情况。
- 捕获DynamoDB请求Payload:如果使用AWS SDK,通过
SdkHttpClient.Builder配置日志拦截器,开启请求日志,记录发送给DynamoDB的完整请求内容,其中包含所有待写入条目,异常发生后可直接从日志中提取重复主键的记录。
二、自定义写入流程的异常处理逻辑
1. 扩展Sink实现自定义异常捕获
如果基于流处理框架(如Flink、Kafka Connect)使用DynamoDB Sink,可以扩展官方Sink类,在写入环节捕获DynamoDbException,针对重复主键错误做针对性处理:
public class CustomDynamoDbBatchSink extends RichSinkFunction<List<YourRecord>> { private DynamoDbAsyncClient ddbClient; private static final Logger LOG = LoggerFactory.getLogger(CustomDynamoDbBatchSink.class); @Override public void open(Configuration params) { ddbClient = DynamoDbAsyncClient.create(); } @Override public void invoke(List<YourRecord> batch, Context context) throws Exception { BatchWriteItemRequest batchRequest = buildBatchWriteRequest(batch); try { BatchWriteItemResponse response = ddbClient.batchWriteItem(batchRequest).get(); // 处理未成功写入的条目 handleUnprocessedItems(response.unprocessedItems(), batch); } catch (ExecutionException e) { if (e.getCause() instanceof DynamoDbException ddbEx) { if (ddbEx.getMessage().contains("Provided list of item keys contains duplicate")) { // 遍历当前批次,找出重复主键的记录 Set<String> seenKeys = new HashSet<>(); for (YourRecord record : batch) { String key = getRecordKey(record); // 自定义方法生成主键字符串 if (!seenKeys.add(key)) { LOG.error("Duplicate key found: {}, record details: {}", key, record); // 可选:将错误记录存入错误队列或存储 saveToErrorStore(record); } } // 可选:跳过错误批次或仅跳过重复项继续处理 return; } throw ddbEx; } throw e; } } private void handleUnprocessedItems(Map<String, List<WriteRequest>> unprocessedItems, List<YourRecord> originalBatch) { if (!unprocessedItems.isEmpty()) { for (List<WriteRequest> writes : unprocessedItems.values()) { for (WriteRequest writeReq : writes) { String key = extractKeyFromWriteRequest(writeReq); LOG.error("Unprocessed item with duplicate key: {}", key); } } } } @Override public void close() { ddbClient.close(); } }
2. 前置校验避免异常触发
在数据进入Sink之前添加预处理步骤,提前校验主键唯一性:
- 流处理场景:使用窗口或状态管理(如Flink的
ValueState)跟踪已处理的主键,过滤掉重复记录; - 批量写入场景:组装写入批次时,用哈希集合存储已出现的主键,实时过滤重复项并记录。
3. 利用框架错误回调机制
多数流处理框架提供内置的错误回调接口,比如:
- Flink中可通过
SinkFunction的invoke方法捕获异常,或使用ProcessingTimeService处理异步错误; - Kafka Connect中可配置
error.tolerance参数,并自定义ErrorReporter将错误记录写入专门的主题,后续再分析这些错误数据。
内容的提问来源于stack exchange,提问作者Thor
相关产品推荐
相关产品推荐

