Kafka Streams 流处理 DSL
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-input和word-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));
}
}
运行与验证
- 启动程序
- 向
word-count-input发送一些消息:“hello kafka streams hello world” - 消费
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 节点:对应各种变换,如
filter、map - Stateful 节点:聚合、连接等会创建状态存储,并在内部变更日志主题中备份
- Sink 节点:通过
to()写回 Kafka
所有操作在逻辑上分解为子拓扑(sub-topology),当遇到 groupBy 或 through 等需要改变分区的操作时,就会生成新的子拓扑,数据会通过 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 - 使用
Producer和Consumer的 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());