代码之家  ›  专栏  ›  技术社区  ›  cscan ssice

Kafka流不会使用聚合值重新启动

  •  0
  • cscan ssice  · 技术社区  · 8 年前

    我在一个流上聚合值,如下所示:

    private KTable<String, StringAggregator> aggregate(KStream<String, String> inputStream) {
        return inputStream
                .groupByKey(Serialized.with(Serdes.String(), Serdes.String()))
                .aggregate(
                        StringAggregator::new,
                        (k, v, a) -> {
                            a.add(v);
                            return a;
                        }, Materialized.<String, StringAggregator>as(Stores.persistentKeyValueStore("STATE_STORE"))
                                .withKeySerde(Serdes.String())
                                .withValueSerde(getValueSerde(StringAggregator.class)));
    }
    

    1 回复  |  直到 8 年前
        1
  •  0
  •   cscan ssice    7 年前

    我最终创建了聚合逻辑,它使用了在卡夫卡主题上持久化的聚合结果。逻辑如下:

    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");