如何在单个Flink管道中基于不同键进行聚合
复用Flink管道实现双Key分组聚合
不需要新建独立管道,直接复用现有Couchbase源的数据流分支处理即可,这样能避免重复读取Couchbase,是更高效的方案。
修改思路
将预处理后的主数据流保存为变量,基于这个变量分别衍生出两个分支,各自完成key1和key2的分组聚合逻辑,最后分别对接Sink即可。
修改后的代码示例
// 初始化并预处理主数据流(仅读取一次Couchbase) DataStream<YourRecordType> mainDataStream = env.addSource(readFromCouchBase...) .name("couchbase-flinkjob") .assignTimestampsAndWatermarks( new TimestampExtractorAndWatermarkEmitter(60 * 1000, false)); // 分支1:基于key1的窗口聚合 DataStream<CountResult> key1AggregationStream = mainDataStream .keyBy(record -> record.getKey1()) .timeWindow(Time.seconds(10)) .aggregate(new CountGroupFunctionWithEventTimeProcessing, new CountGroupWindowFunction); // 分支2:基于key2的窗口聚合 DataStream<CountResult> key2AggregationStream = mainDataStream .keyBy(record -> record.getKey2()) .timeWindow(Time.seconds(10)) .aggregate(new CountGroupFunctionWithEventTimeProcessing, new CountGroupWindowFunction); // 将两个聚合结果分别写入Postgres key1AggregationStream.addSink(new IdempotentPostgresSqlSinkFunction).name("postgres-sink-key1"); key2AggregationStream.addSink(new IdempotentPostgresSqlSinkFunction).name("postgres-sink-key2");
关键说明
- 主数据流仅从Couchbase读取一次,后续所有聚合分支都复用这股数据流,避免了重复拉取数据的资源消耗
- 如果两个聚合结果需要写入同一张Postgres表,可通过
union()方法合并两个流后再写入;若写入不同表,分别配置Sink即可
内容的提问来源于stack exchange,提问作者user3553913
相关产品推荐
相关产品推荐

