我最终创建了聚合逻辑,它使用了在卡夫卡主题上持久化的聚合结果。逻辑如下:
private KStream<String, StringAggregator> getAggregator(String topicName,
KStream<String, String> input,
KTable<String, StringAggregator> aggregator) {
return input
.leftJoin(aggregator, (inputMessage, aggregatorMessage) -> {
if (aggregatorMessage == null) {
aggregatorMessage = new StringAggregator();
}
aggregatorMessage.add(inputMessage);
return aggregatorMessage;
}).peek((k, v) -> logger.info("Aggregated a join input for {}: {}, {} aggregated.", topicName, k, v.size()));
}
下面是实际构建流的逻辑。
String topicName = "input";
KStream<String, String> input = streamsBuilder.stream(topicName);
KTable<String, StringAggregator> aggregator = streamsBuilder.table("aggregate");
getAggregator(topicName, input, aggregator).to("aggregate");