单个聚合结果写入多个DynamoDB条目实现方案问询
Flink DynamoDbSink 实现单聚合结果写入多分区键方案
无需重写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
相关产品推荐
相关产品推荐

