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

单个聚合结果写入多个DynamoDB条目实现方案问询

无需重写DynamoDbSink的前提下,有两种可行方案实现单个聚合结果写入多个分区键条目:

方案一:数据流前置拆分

在DynamoDbSink之前添加flatMap算子,将单个聚合对象拆分为多个对应不同分区键的独立数据条目,再将拆分后的数据流传入Sink:

// 聚合后的数据流
DataStream<AggregatedResult> aggregatedStream = ...;

// 拆分聚合结果为多分区键条目
DataStream<PartitionedItem> partitionedStream = aggregatedStream.flatMap((AggregatedResult value, Collector<PartitionedItem> out) -> {
    // 根据业务需求生成对应不同分区键的条目
    out.collect(new PartitionedItem("partition-key-1", value.getAggregatedData()));
    out.collect(new PartitionedItem("partition-key-2", value.getAggregatedData()));
    // 可按需添加更多分区键条目
});

// 配置DynamoDbSink,每个拆分后的条目对应一个写入请求
DynamoDbSink<PartitionedItem> sink = DynamoDbSink.<PartitionedItem>builder()
        .setDynamoDbClient(dynamoDbClient)
        .setElementConverter(item -> {
            Map<String, AttributeValue> attributes = new HashMap<>();
            attributes.put("partitionKey", AttributeValue.builder().s(item.getPartitionKey()).build());
            attributes.put("aggData", AttributeValue.builder().s(item.getAggData()).build());
            return DynamoDbWriteRequest.builder()
                    .putRequest(PutRequest.builder().item(attributes).build())
                    .build();
        })
        .setTableName("target-table")
        .build();

partitionedStream.addSink(sink);

方案二:使用批量转换接口

利用Flink提供的DynamoDbBatchConverter接口(Flink 1.15+版本支持),直接将单个聚合对象转换为多个DynamoDbWriteRequest,无需拆分数据流:

// 聚合后的数据流
DataStream<AggregatedResult> aggregatedStream = ...;

// 配置DynamoDbSink,通过BatchConverter生成多写入请求
DynamoDbSink<AggregatedResult> sink = DynamoDbSink.<AggregatedResult>builder()
        .setDynamoDbClient(dynamoDbClient)
        .setBatchConverter(value -> {
            List<DynamoDbWriteRequest> requests = new ArrayList<>();
            // 生成第一个分区键的写入请求
            Map<String, AttributeValue> item1 = new HashMap<>();
            item1.put("partitionKey", AttributeValue.builder().s("key1").build());
            item1.put("aggData", AttributeValue.builder().s(value.getAggregatedData()).build());
            requests.add(DynamoDbWriteRequest.builder()
                    .putRequest(PutRequest.builder().item(item1).build())
                    .build());
            
            // 生成第二个分区键的写入请求
            Map<String, AttributeValue> item2 = new HashMap<>();
            item2.put("partitionKey", AttributeValue.builder().s("key2").build());
            item2.put("aggData", AttributeValue.builder().s(value.getAggregatedData()).build());
            requests.add(DynamoDbWriteRequest.builder()
                    .putRequest(PutRequest.builder().item(item2).build())
                    .build());
            
            return requests;
        })
        .setTableName("target-table")
        .build();

aggregatedStream.addSink(sink);

方案对比

  • 方案一逻辑直观,数据流拆分后Sink的转换逻辑简单,便于后续维护和调试。
  • 方案二更精简,无需额外算子,直接在转换层完成多请求生成,适合对数据流链路有精简需求的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 03:07:39