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

如何让Kafka Streams仅输出groupBy最终聚合结果至Kafka Topic

Kafka Streams 仅输出最终聚合结果的实现方案

现有Kafka Streams代码通过groupByKey()和reduce()对输入Topic中的Person数据进行聚合,当前输出Topic会存储count字段的所有更新历史。需求为:仅在处理完输入Topic(数据来自CSV文件,单行为一条消息,总量达数百万条)的所有消息后,将groupBy的最终聚合结果输出至输出Topic。

可行方案

方案1:基于结束信号触发输出

在CSV文件全部导入到输入Topic后,向该Topic发送一条特殊的"结束信号"消息(比如key设为__END__)。Streams应用监听该信号,收到后遍历聚合状态存储,将所有最终结果发送到输出Topic。

修改后的流处理代码

public static void createGroupByStream(final StreamsBuilder builder) {
    // 定义聚合状态存储的名称
    final String AGG_STORE_NAME = "person-aggregation-store";
    
    // 读取输入流
    KStream<String, Person> inputStream = builder.stream(
        "kafka-stream-input",
        Consumed.with(Serdes.String(), CustomSerdes.PersonSerde())
    );

    // 1. 执行聚合逻辑,将结果存入状态存储(不直接输出到Topic)
    inputStream
        // 过滤掉结束信号,避免干扰聚合
        .filter((key, person) -> !"__END__".equals(key))
        .groupByKey()
        .reduce((prev, curr) -> {
            // 注意:创建新对象避免修改状态存储中的原实例,防止并发问题
            Person mergedPerson = new Person();
            mergedPerson.setName(prev.getName());
            mergedPerson.setCount(prev.getCount() + curr.getCount());
            return mergedPerson;
        }, Materialized.as(AGG_STORE_NAME)) // 显式指定状态存储
        .toStream()
        .filter((key, value) -> false); // 临时过滤所有中间输出

    // 2. 监听结束信号,触发最终结果输出
    inputStream
        .filter((key, person) -> "__END__".equals(key))
        .foreach((key, signal) -> {
            // 获取状态存储的只读实例
            ReadOnlyKeyValueStore<String, Person> aggStore = 
                KafkaStreamsUtils.getReadOnlyKeyValueStore(AGG_STORE_NAME);
            
            // 遍历所有聚合结果并发送到输出Topic
            try (KeyValueIterator<String, Person> iterator = aggStore.all()) {
                while (iterator.hasNext()) {
                    KeyValue<String, Person> entry = iterator.next();
                    KafkaProducerUtils.getProducer().send(
                        new ProducerRecord<>(
                            "kafka-stream-output",
                            entry.key,
                            entry.value
                        )
                    );
                }
            }
        });
}

配套工具类示例

需要实现获取状态存储和生产者的工具类:

// KafkaStreamsUtils.java
public class KafkaStreamsUtils {
    private static KafkaStreams streams;

    public static void setKafkaStreams(KafkaStreams streams) {
        KafkaStreamsUtils.streams = streams;
    }

    public static ReadOnlyKeyValueStore<String, Person> getReadOnlyKeyValueStore(String storeName) {
        return streams.store(
            StoreQueryParameters.fromNameAndType(
                storeName,
                QueryableStoreTypes.keyValueStore()
            )
        );
    }
}

// KafkaProducerUtils.java
public class KafkaProducerUtils {
    private static Producer<String, Person> producer;

    static {
        // 配置生产者参数
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class);
        producer = new KafkaProducer<>(props);
    }

    public static Producer<String, Person> getProducer() {
        return producer;
    }
}

使用说明

CSV文件全部导入到输入Topic后,执行以下命令发送结束信号:

kafka-console-producer.sh --broker-list localhost:9092 --topic kafka-stream-input --property parse.key=true --property key.separator=:

然后输入:__END__:(key为__END__,value为空)


方案2:基于交互式查询(Interactive Queries)触发输出

利用Kafka Streams的交互式查询能力,在外部程序(比如REST接口)中查询聚合状态存储,当确认输入Topic已消费完毕时,将结果发送到输出Topic。

实现步骤

  1. 在聚合时显式指定状态存储(同方案1的聚合逻辑)
  2. 启动Streams应用后,通过外部程序查询状态存储:
// 示例:REST接口触发输出
@RestController
public class AggregationExportController {
    private final KafkaStreams streams;
    private final Producer<String, Person> producer;

    public AggregationExportController(KafkaStreams streams, Producer<String, Person> producer) {
        this.streams = streams;
        this.producer = producer;
    }

    @PostMapping("/export-final-results")
    public ResponseEntity<Void> exportResults() {
        // 1. 先验证输入Topic是否已消费完毕(通过查询consumer group偏移量)
        if (!isInputTopicConsumedCompletely()) {
            return ResponseEntity.status(HttpStatus.PRECONDITION_FAILED).build();
        }

        // 2. 查询聚合状态存储
        ReadOnlyKeyValueStore<String, Person> aggStore = streams.store(
            StoreQueryParameters.fromNameAndType(
                "person-aggregation-store",
                QueryableStoreTypes.keyValueStore()
            )
        );

        // 3. 发送结果到输出Topic
        try (KeyValueIterator<String, Person> iterator = aggStore.all()) {
            while (iterator.hasNext()) {
                KeyValue<String, Person> entry = iterator.next();
                producer.send(new ProducerRecord<>("kafka-stream-output", entry.key, entry.value));
            }
        }

        return ResponseEntity.ok().build();
    }

    // 实现输入Topic消费完成的校验逻辑
    private boolean isInputTopicConsumedCompletely() {
        // 逻辑:查询consumer group的当前偏移量与Topic的最新偏移量是否相等
        // 可通过AdminClient实现,此处省略具体代码
        return true;
    }
}

方案3:基于预设总消息数触发输出

如果能提前统计CSV文件的总行数(即输入Topic的总消息数),可以在Streams应用中统计消费的消息总数,达到预设值时自动输出最终结果。

修改后的流处理代码

public static void createGroupByStream(final StreamsBuilder builder) {
    final String AGG_STORE_NAME = "person-aggregation-store";
    final String COUNT_STORE_NAME = "message-count-store";
    // 从配置文件读取预设的总消息数(CSV总行数)
    final long TOTAL_MESSAGE_COUNT = 1000000;

    KStream<String, Person> inputStream = builder.stream(
        "kafka-stream-input",
        Consumed.with(Serdes.String(), CustomSerdes.PersonSerde())
    );

    // 1. 统计消费的总消息数
    inputStream
        .groupBy((key, value) -> "total") // 用固定key统计全局总数
        .count(Materialized.as(COUNT_STORE_NAME))
        .toStream()
        .foreach((key, currentCount) -> {
            if (currentCount >= TOTAL_MESSAGE_COUNT) {
                // 触发最终结果输出
                ReadOnlyKeyValueStore<String, Person> aggStore = 
                    KafkaStreamsUtils.getReadOnlyKeyValueStore(AGG_STORE_NAME);
                
                try (KeyValueIterator<String, Person> iterator = aggStore.all()) {
                    while (iterator.hasNext()) {
                        KeyValue<String, Person> entry = iterator.next();
                        KafkaProducerUtils.getProducer().send(
                            new ProducerRecord<>("kafka-stream-output", entry.key, entry.value)
                        );
                    }
                }
            }
        });

    // 2. 执行聚合逻辑
    inputStream
        .groupByKey()
        .reduce((prev, curr) -> {
            Person mergedPerson = new Person();
            mergedPerson.setName(prev.getName());
            mergedPerson.setCount(prev.getCount() + curr.getCount());
            return mergedPerson;
        }, Materialized.as(AGG_STORE_NAME));
}

注意事项

  • 聚合逻辑中不要修改原Person对象,必须创建新对象返回,否则会导致状态存储中的数据被意外修改,引发并发问题。
  • 状态存储建议使用持久化存储(比如RocksDB),避免应用重启后丢失聚合数据。
  • 方案1的结束信号必须在CSV全部导入后发送,否则会提前输出不完整的结果。
  • 方案3的总消息数必须准确,否则会出现提前输出或永远不输出的情况。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 19:25:02