Kafka Streams 流处理 DSL

FreeGuideOnline 最新 2026-07-09

java KStream<String, String> stream = builder.stream("input-topic");


### KTable(变更日志表)

**KTable** 代表一个**可更新的数据集**,类似数据库表。每个键对应的值可以被后续消息覆盖(upsert),因此它始终反映每个键的“最新状态”。

- 同样基于键值对,但后续相同键的消息会代替之前的值
- 删除操作通过发送 `null` 值(tombstone)实现
- 适合处理“事实表”,如用户资料、账户余额

```java
KTable<String, String> table = builder.table("user-info-topic");

GlobalKTable

GlobalKTable 是 KTable 的特殊形式,它将整个数据集复制到每个 Kafka Streams 实例中,用于与 KStream 进行无键关联的查找连接。它通常用于小规模、变化不频繁的维度表。

直观对比

特性 KStream KTable
数据模型 事件流 可更新表
键重复时 保留所有记录 只保留最新记录
查询上下文 无法直接查询(除IQ) 可查询状态存储
典型场景 点击流、交易事件 产品编目、用户最新资料

DSL 基础操作

Kafka Streams DSL 提供了丰富的内置操作,它们都返回新的 KStream/KTable,因此可以链式调用构建完整拓扑。

无状态转换

这些操作每次只处理一条记录,不依赖状态。

map / mapValues

对记录的值或键值进行一对一映射。

// 将值转为大写
stream.mapValues(value -> value.toUpperCase());

// 同时更改键和值
stream.map((key, value) -> KeyValue.pair("prefix_" + key, value));

filter / filterNot

根据条件保留或丢弃记录。

stream.filter((key, value) -> value.length() > 10);

flatMap / flatMapValues

一对多变换,一条记录可以变成零到多条记录。

stream.flatMapValues(value -> Arrays.asList(value.split("\\s+")));

through

将流写入一个中间主题并继续处理,常用于调试或强制重分区。

stream.through("intermediate-topic");

有状态操作

需要保存和查询状态,以处理聚合、连接等。

groupBy / groupByKey

将记录分组,为后续聚合做准备。groupBy 会引发重分区(数据重新分布到不同任务),groupByKey 仅当键未变时保持分区。

KGroupedStream<String, String> grouped = stream
    .groupBy((key, value) -> value.substring(0, 3)); // 按值的前三位分组

聚合(aggregate, count, reduce)

在分组后统计、求和、归并等。

// 计数
KTable<String, Long> counts = grouped.count();

// 求和(需要初始值)
KTable<String, Integer> sum = grouped.aggregate(
    () -> 0,
    (aggKey, newValue, aggValue) -> aggValue + Integer.parseInt(newValue),
    Materialized.with(Serdes.String(), Serdes.Integer())
);

注意:聚合结果存储在本地状态仓库(RocksDB),并可选地写回 Kafka 变更日志主题以实现容错。

连接(Join)

流与流、流与表都可以连接。

  • KStream-KStream Join:基于时间窗口,仅在窗口期内的匹配记录会输出。
  • KStream-KTable Join:每次 KStream 来一条记录,查询对应 KTable 中该键的最新状态,非窗口操作。
  • KStream-GlobalKTable Join:全量维度表关联,不按键分区。
KStream<String, String> joined = stream.join(table,
    (streamValue, tableValue) -> streamValue + " - " + tableValue,
    Joined.with(Serdes.String(), Serdes.String(), Serdes.String()));

窗口化

聚合时常需要时间窗口。Kafka Streams 内置多种窗口:

  • Tumbling time windows(翻滚窗口): 固定长度,不重叠
  • Hopping time windows(滑动窗口): 固定长度,有重叠步长
  • Session windows(会话窗口): 基于活动间隙动态划分
  • Sliding windows(仅用于连接): 固定大小的滑动窗口
TimeWindows windowDef = TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(5));
KTable<Windowed<String>, Long> windowedCount = grouped
    .windowedBy(windowDef)
    .count();

窗口操作通常需要指定宽限期(grace period) 以处理迟到数据,并可以通过 suppress 抑制中间结果。


构建第一个 Kafka Streams DSL 应用

下面是一个完整的入门示例:统计一段文本中单词的出现次数,实时输出到控制台。

环境准备

  • Apache Kafka 已启动(默认 localhost:9092)
  • 创建两个主题:word-count-inputword-count-output
  • Maven 依赖:
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams</artifactId>
    <version>3.7.0</version>
</dependency>
<dependency>
    <groupId>org.slf4j</groupId>
    <artifactId>slf4j-simple</artifactId>
    <version>2.0.9</version>
</dependency>

编写代码

import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.*;
import org.apache.kafka.streams.kstream.*;

import java.util.Arrays;
import java.util.Properties;

public class WordCountSample {
    public static void main(String[] args) {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-application");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        StreamsBuilder builder = new StreamsBuilder();

        KStream<String, String> textLines = builder.stream("word-count-input");

        KTable<String, Long> wordCounts = textLines
                .flatMapValues(text -> Arrays.asList(text.toLowerCase().split("\\W+")))
                .selectKey((key, word) -> word)                // 变更键为单词本身
                .groupByKey()                                   // 按单词分组
                .count();                                       // 计数

        wordCounts.toStream().to("word-count-output", Produced.with(Serdes.String(), Serdes.Long()));

        KafkaStreams streams = new KafkaStreams(builder.build(), props);
        streams.start();

        Runtime.getRuntime().addShutdownHook(new Thread(streams::close));
    }
}

运行与验证

  1. 启动程序
  2. word-count-input 发送一些消息:“hello kafka streams hello world”
  3. 消费 word-count-output,会看到:
hello   2
kafka   1
streams 1
world   1

这就是一个基于 DSL 的经典 Word Count 流处理程序。


理解 DSL 计算的底层机制

Kafka Streams 的 DSL 会自动创建拓扑(Topology)。每个操作都被翻译成拓扑节点,背后是由 Processor API 实现。了解这些细节能帮你排错和优化。

  • Source 节点:对应 builder.stream()builder.table()
  • Processor 节点:对应各种变换,如 filtermap
  • Stateful 节点:聚合、连接等会创建状态存储,并在内部变更日志主题中备份
  • Sink 节点:通过 to() 写回 Kafka

所有操作在逻辑上分解为子拓扑(sub-topology),当遇到 groupBythrough 等需要改变分区的操作时,就会生成新的子拓扑,数据会通过 Kafka 中间主题在不同子拓扑间传输。


常用模式与最佳实践

动态路由(Branch)

KStream<String, String>[] branches = stream.branch(
    (key, value) -> value.contains("error"),
    (key, value) -> value.contains("warn"),
    (key, value) -> true  // 其他
);
branches[0].to("error-topic");
branches[1].to("warn-topic");
branches[2].to("info-topic");

表-表连接(Table-Table Join)

虽然 KTable 本身可更新,但它也支持与其他 KTable 做连接(非窗口),常用于实现物化视图。

KTable<String, String> joinedTable = table1.join(table2,
    (value1, value2) -> value1 + "," + value2);

使用状态存储进行业务逻辑

有时需要更细粒度的控制,可以通过 transformValues 配合 ValueTransformer 访问状态存储,实现自定义逻辑。

stream.transformValues(() -> new MyTransformer(), "my-state-store");

可靠性与容错

  • 始终设置唯一的 APPLICATION_ID_CONFIG
  • 使用 ProducerConsumer 的 EOS(精确一次语义):processing.guarantee=exactly_once_v2
  • 状态存储会自动基于变更日志恢复

测试 DSL 拓扑

Kafka Streams 提供了 TopologyTestDriver,让你可以脱离 Kafka 集群进行单元测试:

TopologyTestDriver testDriver = new TopologyTestDriver(builder.build(), props);
TestInputTopic<String, String> input = testDriver.createInputTopic("word-count-input",
    Serdes.String().serializer(), Serdes.String().serializer());
TestOutputTopic<String, Long> output = testDriver.createOutputTopic("word-count-output",
    Serdes.String().deserializer(), Serdes.Long().deserializer());

input.pipeInput("hello hello");
assertEquals(2L, output.readValue());