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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 16:02:37