📌 本文是「Kafka 3.9 教程」系列第 八 篇 · 系列目录
本文档基于 Apache Kafka 3.9 官方文档整理
目录
8.1 简介
8.1.1 什么是 Kafka Streams?
Kafka Streams 是 Kafka 内置的轻量级流处理库,用于构建实时应用程序和微服务。
核心特点:
| 特性 | 说明 |
|---|---|
| 轻量级 | 无需独立集群,作为库嵌入应用 |
| 高扩展 | 水平扩展,自动负载均衡 |
| 弹性 | 自动故障转移 |
| 容错 | 支持 Exactly-Once 语义 |
| 易用 | Java/Scala API,简洁优雅 |
8.1.2 与其他流处理框架对比
| 特性 | Kafka Streams | Apache Flink | Spark Streaming |
|---|---|---|---|
| 部署 | 嵌入应用 | 独立集群 | 独立集群 |
| 延迟 | 毫秒级 | 毫秒级 | 秒级 |
| 处理模型 | 纯流处理 | 纯流处理 | 微批处理 |
| 状态管理 | 内置 RocksDB | 内置 | 需额外配置 |
| Exactly-Once | 支持 | 支持 | 支持(Spark 3.0+) |
| 学习曲线 | 低 | 中 | 中 |
8.1.3 应用场景
- 实时数据处理:ETL、转换、聚合
- 事件驱动应用:告警、通知、推荐
- 复杂事件处理:模式匹配、CEP
- 实时分析:指标计算、报表
8.2 核心概念
8.2.1 KStream
KStream 表示无界事件流,类似于数据库表。
特点:
- 每条消息都是独立的
- 可以重复消费
- 适用于原始事件数据
KStream 示意:
┌──────────────────────────────────────────────────────────┐
│ KStream<String, Long> │
│ │
│ ("apple", 1) ──▶ ("banana", 1) ──▶ ("apple", 1) │
│ │
│ 每条记录独立,没有关联 │
└──────────────────────────────────────────────────────────┘
8.2.2 KTable
KTable 表示变更日志流,类似于数据库视图。
特点:
- 相同 Key 的记录会合并
- 始终保持最新状态
- 适用于聚合结果、维表
KTable 示意:
┌──────────────────────────────────────────────────────────┐
│ KTable<String, Long> (word count) │
│ │
│ ("apple", 5) ←── ("apple", +1) 更新 │
│ │
│ 相同 Key 的记录会合并 │
└──────────────────────────────────────────────────────────┘
8.2.3 GlobalKTable
GlobalKTable 是全局 KTable。
特点:
- 每个实例都有完整数据
- 不参与分区
- 适用于维表关联
KStream + GlobalKTable 关联示意:
┌─────────────────┐ ┌─────────────────┐
│ Orders (KStream)│ │ Products (GKT) │
│ │ │ │
│ product_id: 1 │──┐ │ id: 1 → Apple │
│ product_id: 2 │ │ │ id: 2 → Banana │
└─────────────────┘ │ └─────────────────┘
│
▼ join
┌─────────────────┐
│ Enriched Order │
│ (product_name) │
└─────────────────┘
8.2.4 State Store
状态存储是有状态处理的基础,用于存储和查询处理过程中的中间结果。
状态存储类型
| 类型 | 说明 | 适用场景 |
|---|---|---|
| RocksDB | 默认的持久化状态存储,基于 RocksDB 嵌入式数据库 | 生产环境,大数据量 |
| 内存 | 内存状态存储,基于 HashMap | 测试,小数据量 |
| 持久化 | 支持 changelog 主题持久化 | 容错恢复 |
状态存储类型详解
// RocksDB 状态存储(默认)
Materialized.as("store-name")
.withLoggingEnabled(true) // 启用 changelog
.withCachingEnabled(true); // 启用缓存
// 内存状态存储
Materialized.as(
Stores.inMemoryKeyValueStore("store-name")
).withCachingEnabled(true);
状态存储操作
import org.apache.kafka.streams.state.QueryableStoreType;
import org.apache.kafka.streams.state.ReadOnlyKeyValueStore;
// 通过交互查询访问状态
ReadOnlyKeyValueStore<String, Long> store = streams.store(
StoreQueryParameters.fromNameAndType(
"word-count-store",
QueryableStoreType.keyValueStore()
)
);
// 查询操作
Long count = store.get("kafka");
// 范围查询
KeyValueIterator<String, Long> range = store.range("a", "z");
while (range.hasNext()) {
KeyValue<String, Long> kv = range.next();
System.out.println(kv.key + " = " + kv.value);
}
状态存储类型
| 存储类型 | 接口 | 说明 |
|---|---|---|
| KeyValueStore | KeyValueStore | 键值存储 |
| WindowStore | WindowStore | 窗口键值存储 |
| SessionStore | SessionStore | 会话存储 |
状态容错
# 启用 changelog(默认启用)
default.kafka.streams.state.store.logging.enabled=true
# changelog 主题配置
default.kafka.streams.state.store.changelog.enable=true
# 复制因子(默认与内部主题相同)
default.kafka.streams.state.store.replication.factor=3
状态缓存
// 配置缓存
Materialized.as("store-name")
.withCachingEnabled(true) // 默认启用
.withCacheMaxBytes(1024 * 1024); // 缓存大小
8.2.5 Topology
拓扑是流处理应用的逻辑结构。
Topology 示意:
┌─────────────────────────────────────────────────────────┐
│ Topology │
│ │
│ ┌─────────┐ ┌─────────┐ ┌─────────┐ │
│ │ Source │───▶│ Process │───▶│ Sink │ │
│ │ (输入) │ │ (处理) │ │ (输出) │ │
│ └─────────┘ └─────────┘ └─────────┘ │
│ │ │
│ ▼ │
│ ┌─────────────┐ │
│ │ State Store │ │
│ └─────────────┘ │
└─────────────────────────────────────────────────────────┘
8.3 DSL API
8.3.1 基本使用
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.kstream.KTable;
import org.apache.kafka.streams.kstream.Produced;
import java.util.Arrays;
import java.util.Properties;
public class WordCountExample {
public static void main(String[] args) {
// 配置
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-app");
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();
// 1. 从 Kafka 读取
KStream<String, String> textLines = builder.stream("text-input");
// 2. 处理
KTable<String, Long> wordCounts = textLines
// 拆分每行为单词
.flatMapValues(line -> Arrays.asList(line.toLowerCase().split("\\W+")))
// 按单词分组
.groupBy((key, word) -> word)
// 统计计数
.count();
// 3. 输出到 Kafka
wordCounts.toStream().to("wordcount-output",
Produced.with(Serdes.String(), Serdes.Long()));
// 启动
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
}
}
8.3.2 无状态操作
转换操作
| 操作 | 说明 | 输入 | 输出 |
|---|---|---|---|
| map | 转换键值对 | KStream | KStream |
| mapValues | 只转换值 | KStream/KTable | 相同 |
| flatMap | 展开转换 | KStream | KStream |
| flatMapValues | 展开值 | KStream | KStream |
| filter | 过滤 | KStream/KTable | 相同 |
| selectKey | 重新指定 Key | KStream | KStream |
| merge | 合并流 | KStream | KStream |
// map 示例 - 同时转换键值对
KStream<String, Integer> numbers = source.map(
(key, value) -> KeyValue.pair(key.toUpperCase(), Integer.parseInt(value))
);
// mapValues 示例 - 只转换值
KStream<String, Integer> counts = source.mapValues(value -> value.length());
// flatMap 示例 - 展开为多条记录
KStream<String, String> sentences = words.flatMap(
(key, word) -> Arrays.asList(
KeyValue.pair(key, word.toUpperCase()),
KeyValue.pair(key, word.toLowerCase())
)
);
// flatMapValues 示例 - 展开值为多条记录
KStream<String, String> characters = source.flatMapValues(
value -> Arrays.asList(value.split(""))
);
// filter 示例 - 过滤记录
KStream<String, String> filtered = source.filter(
(key, value) -> value.startsWith("valid")
);
// filterNot 示例 - 过滤不符合条件的记录
KStream<String, String> valid = source.filterNot(
(key, value) -> value.contains("error")
);
// selectKey 示例 - 重新指定键
KStream<String, String> byWord = source.selectKey(
(key, value) -> value.split(" ")[0] // 用第一个单词作为键
);
// merge 示例 - 合并两个流
KStream<String, String> merged = stream1.merge(stream2);
分支操作
// 使用 branch 进行条件分流
KStream<String, String>[] branches = source.branch(
(key, value) -> value.startsWith("priority-1"), // 分支1
(key, value) -> value.startsWith("priority-2"), // 分支2
(key, value) -> true // 默认分支
);
branches[0].to("priority-1-topic");
branches[1].to("priority-2-topic");
branches[2].to("other-topic");
8.3.3 有状态操作
分组操作
| 操作 | 说明 | 输入 | 输出 |
|---|---|---|---|
| groupBy | 按键分组 | KStream | KGroupedStream |
| groupByKey | 按现有键分组 | KStream | KGroupedStream |
// groupBy 示例 - 按指定字段分组
KGroupedStream<String, String> grouped = source.groupBy(
(key, value) -> value.split(":")[0] // 按值的某部分分组
);
// groupByKey 示例 - 按现有键分组
KGroupedStream<String, String> groupedByKey = source.groupByKey();
聚合操作
| 操作 | 说明 | 初始值 | 更新函数 |
|---|---|---|---|
| count | 计数 | 0 | +1 |
| reduce | 归约 | 无 | (oldValue, newValue) -> result |
| aggregate | 聚合 | 自定义 | (key, newValue, aggValue) -> result |
// count 示例
KTable<String, Long> counts = source
.groupBy((key, value) -> value)
.count();
// count 示例 - 使用 Materialized
KTable<String, Long> counts = source
.groupBy((key, value) -> value)
.count(
Materialized.as("word-count-store") // 命名状态存储
.withValueSerde(Serdes.Long())
);
// reduce 示例 - 选择最新值
KTable<String, Event> latestEvent = source
.groupBy((key, value) -> value.getId())
.reduce(
(oldValue, newValue) ->
oldValue.getTimestamp() > newValue.getTimestamp() ? oldValue : newValue,
Materialized.as("latest-event-store")
);
// aggregate 示例 - 求和聚合
KTable<String, Integer> sumAmount = source
.groupBy((key, value) -> value.getCategory())
.aggregate(
() -> 0, // 初始值
(key, newValue, aggValue) -> aggValue + newValue.getAmount(),
Materialized.as("sum-store")
.withValueSerde(Serdes.Integer())
);
8.3.4 表操作
KStream → KTable 转换
// 将 KStream 聚合为 KTable
KTable<String, Long> wordCounts = textLines
.flatMapValues(line -> Arrays.asList(line.toLowerCase().split("\\W+")))
.groupBy((key, word) -> word)
.count();
KTable 操作
// filter - 过滤 KTable
KTable<String, Long> filtered = table.filter(
(key, value) -> value > 100
);
// mapValues - 转换值
KTable<String, Double> normalized = table.mapValues(
value -> value / 1000.0
);
// join - 表连接
KTable<String, String> enriched = leftTable.join(
rightTable,
(leftValue, rightValue) -> leftValue + ":" + rightValue
);
//leftJoin - 左连接
KTable<String, String> enriched = leftTable.leftJoin(
rightTable,
(leftValue, rightValue) -> rightValue != null
? leftValue + ":" + rightValue
: leftValue + ":null"
);
// outerJoin - 外连接
KTable<String, String> enriched = leftTable.outerJoin(
rightTable,
(leftValue, rightValue) -> {
if (leftValue == null) return "null:" + rightValue;
if (rightValue == null) return leftValue + ":null";
return leftValue + ":" + rightValue;
}
);
8.3.5 连接操作
KStream-KStream 连接
// 内连接
KStream<String, String> joined = left.join(
right,
(leftValue, rightValue) -> leftValue + "-" + rightValue,
JoinWindows.of(Duration.ofMinutes(5))
);
// 左连接
KStream<String, String> joined = left.leftJoin(
right,
(leftValue, rightValue) -> rightValue != null
? leftValue + "-" + rightValue
: leftValue + "-null",
JoinWindows.of(Duration.ofMinutes(5))
);
// 外连接
KStream<String, String> joined = left.outerJoin(
right,
(leftValue, rightValue) -> {
if (leftValue == null && rightValue == null) return "null-null";
if (leftValue == null) return "null-" + rightValue;
if (rightValue == null) return leftValue + "-null";
return leftValue + "-" + rightValue;
},
JoinWindows.of(Duration.ofMinutes(5))
);
KStream-KTable 连接
// 内连接(需要协同分区)
KStream<String, EnrichedOrder> enriched = orders.join(
products,
(orderId, order) -> order.getProductId(),
(order, product) -> new EnrichedOrder(order, product)
);
// 左连接
KStream<String, EnrichedOrder> enriched = orders.leftJoin(
products,
(orderId, order) -> order.getProductId(),
(order, product) -> {
if (product == null) {
return new EnrichedOrder(order, null);
}
return new EnrichedOrder(order, product);
}
);
KStream-GlobalKTable 连接
// GlobalKTable 连接(无需协同分区)
GlobalKTable<String, Product> globalProducts = builder.globalTable(
"products-topic",
Consumed.with(Serdes.String(), productSerde),
Materialized.as("global-products-store")
);
KStream<String, EnrichedOrder> enriched = orders.join(
globalProducts,
(order, productId) -> order.getProductId(),
(order, product) -> new EnrichedOrder(order, product)
);
8.3.6 输出操作
发送到主题
// to - 发送值
wordCounts.toStream().to("wordcount-output");
// to - 指定序列化器
wordCounts.toStream().to(
"wordcount-output",
Produced.with(Serdes.String(), Serdes.Long())
);
转换为 KStream
// toStream - 将 KTable 转换为 KStream
KStream<String, Long> wordStream = wordCounts.toStream();
// toStream - 只发送更新
wordCounts.toStream().to("wordcount-updates");
// toStream - 发送更改记录
wordCounts.toStream(
// 发送 INSERT, UPDATE, DELETE 等更改记录
Produced.with(Serdes.String(), Serdes.Long())
).to("wordcount-changelog");
聚合操作
| 操作 | 说明 | 类型 |
|---|---|---|
| count | 计数 | KTable |
| reduce | 归约 | KTable |
| aggregate | 聚合 | KTable |
// count 示例
KTable<String, Long> counts = source
.groupBy((key, value) -> value)
.count();
// reduce 示例
KTable<String, String> concatenated = source
.groupBy((key, value) -> KeyValue.pair(value.getCategory(), value))
.reduce((oldValue, newValue) ->
oldValue.getCount() > newValue.getCount() ? oldValue : newValue
);
// aggregate 示例
KTable<String, Integer> summed = source
.groupBy((key, value) -> KeyValue.pair(value.getCategory(), value.getAmount()))
.aggregate(
() -> 0, // 初始值
(key, newValue, aggValue) -> aggValue + newValue,
Materialized.as("sum-store")
);
连接操作
| 操作 | 说明 | 要求 |
|---|---|---|
| join | 内连接 | 相同时间语义 |
| leftJoin | 左连接 | 协同分区 |
| outerJoin | 外连接 | 相同 Key 类型 |
// KStream - KStream Join
KStream<String, OrderEnriched> enrichedOrders = orders.join(
products,
(order, product) -> new OrderEnriched(order, product),
JoinWindows.of(Duration.ofMinutes(5)),
Joined.with(Serdes.String(), orderSerde, productSerde)
);
// KStream - KTable Join (维表关联)
KStream<String, OrderWithProduct> orderWithProduct = orders.join(
productTable,
(orderId, order) -> order.getProductId(),
(order, product) -> new OrderWithProduct(order, product)
);
// KStream - GlobalKTable Join (无需协同分区)
KStream<String, OrderWithProduct> orderWithProduct = orders.join(
globalProductTable,
(order, productId) -> productId,
(order, product) -> new OrderWithProduct(order, product)
);
8.3.7 Serde(序列化/反序列化)
Kafka Streams 需要为所有数据类型配置序列化器/反序列化器。
内置 Serde
import org.apache.kafka.common.serialization.Serdes;
// 基本类型 Serde
Serdes.String() // String
Serdes.Integer() // Integer
Serdes.Long() // Long
Serdes.Double() // Double
Serdes.Float() // Float
Serdes.Boolean() // Boolean
Serdes.ByteBuffer() // ByteBuffer
Serdes.Bytes() // byte[]
// List, Set, Map 等集合类型
JSON Serde
// 使用 JsonSerde
import org.apache.kafka.connect.json.JsonSerde;
// 配置 JSON Serde
JsonSerde<MyRecord> myRecordSerde = new JsonSerde<>(MyRecord.class);
// 在拓扑中使用
streamsBuilder.stream("input-topic",
Consumed.with(Serdes.String(), myRecordSerde)
);
Avro Serde
// 使用 Confluent Schema Registry
import io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde;
Properties props = new Properties();
props.put("schema.registry.url", "http://localhost:8081");
props.put("basic.auth.credentials.source", "USER_INFO");
props.put("basic.auth.user.info", "user:password");
// Avro Serde 配置
Serde<MyAvroRecord> avroSerde = new SpecificAvroSerde<>();
自定义 Serde
import org.apache.kafka.common.serialization.Serializer;
import org.apache.kafka.common.serialization.Deserializer;
import org.apache.kafka.common.serialization.Serde;
// 自定义序列化器
public class JsonSerializer implements Serializer<MyRecord> {
@Override
public byte[] serialize(String topic, MyRecord data) {
// JSON 序列化
return gson.toJson(data).getBytes(StandardCharsets.UTF_8);
}
}
// 自定义反序列化器
public class JsonDeserializer implements Deserializer<MyRecord> {
@Override
public MyRecord deserialize(String topic, byte[] data) {
// JSON 反序列化
return gson.fromJson(new String(data, StandardCharsets.UTF_8), MyRecord.class);
}
}
// 自定义 Serde
public class JsonSerde implements Serde<MyRecord> {
@Override
public Serializer<MyRecord> serializer() {
return new JsonSerializer();
}
@Override
public Deserializer<MyRecord> deserializer() {
return new JsonDeserializer();
}
}
配置默认 Serde
Properties props = new Properties();
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.Long().getClass());
8.3.8 时间语义
Kafka Streams 支持不同的时间提取和处理策略。
时间提取配置
# 默认时间戳提取器
default.kafka.streams.timestamp.extractor=org.apache.kafka.streams.processor.FailOnInvalidTimestamp
# 其他选项
# LogAppendTime - 使用 Broker 时间
# WallclockTimestampExtractor - 使用系统时间
# MetadataTimestampExtractor - 从元数据提取时间
自定义时间戳提取器
import org.apache.kafka.streams.processor.TimestampExtractor;
public class MyTimestampExtractor implements TimestampExtractor {
@Override
public long extract(ConsumerRecord<Object, Object> record, long partitionTime) {
// 从记录值中提取时间戳
String value = (String) record.value();
return Long.parseLong(value.split(",")[0]); // 假设第一个字段是时间戳
}
}
窗口时间语义
// 流处理时间
TimeWindows.ofSizeWithNoGrace(Duration.ofMinutes(1)) // 不使用延迟数据
// 事件时间(支持延迟数据)
TimeWindows.ofSizeAndGrace(
Duration.ofMinutes(1), // 窗口大小
Duration.ofSeconds(30) // 优雅期(允许延迟数据)
)
// 会话窗口(不连续事件)
SessionWindows.with(Duration.ofMinutes(5))
8.3.9 命名主题与资源管理
自定义内部主题名称
# 禁用自动创建
# spring.kafka.streams.auto.create.topics.enable=false
# 命名规则
# <application.id>-<operatorName>-<changelog/partition/repartition>
配置内部主题
# 复制因子
replication.factor=3
# 分区数
num.stream.threads=3
# 最小同步副本数
min.insync.replicas=2
8.4 Processor API
8.4.1 基本用法
import org.apache.kafka.streams.processor.Processor;
import org.apache.kafka.streams.processor.ProcessorContext;
import org.apache.kafka.streams.processor.ProcessorSupplier;
import org.apache.kafka.streams.Topology;
public class MyProcessorSupplier implements ProcessorSupplier<String, String> {
@Override
public Processor<String, String> get() {
return new Processor<String, String>() {
private ProcessorContext context;
@Override
public void init(ProcessorContext context) {
this.context = context;
this.context.schedule(Duration.ofSeconds(1));
}
@Override
public void process(String key, String value) {
// 处理消息
System.out.println("Received: " + key + " = " + value);
// 转发到下游
context.forward(key, value.toUpperCase());
// 提交偏移量
context.commit();
}
@Override
public void close() {
// 清理资源
}
};
}
}
8.4.2 构建拓扑
import org.apache.kafka.streams.Topology;
Topology topology = new Topology();
// 添加处理器
topology.addSource("source", "input-topic")
.addProcessor("processor", () -> new MyProcessorSupplier(), "source")
.addSink("sink", "output-topic", "processor");
8.5 窗口操作
8.5.1 窗口类型
| 窗口类型 | 说明 | 使用场景 |
|---|---|---|
| Tumbling | 滚动窗口 | 每分钟统计 |
| Hopping | 跳跃窗口 | 每分钟滑动30秒 |
| Sliding | 滑动窗口 | 实时移动平均 |
| Session | 会话窗口 | 用户会话分析 |
8.5.2 滚动窗口(Tumbling)
import org.apache.kafka.streams.kstream.TimeWindows;
import org.apache.kafka.streams.kstream.Windowed;
TimeWindows window = TimeWindows.ofSizeAndGrace(
Duration.ofMinutes(1), // 窗口大小
Duration.ofSeconds(30) // 优雅期
);
KTable<Windowed<String>, Long> windowedCounts = textLines
.groupBy((key, value) -> value)
.windowedBy(window)
.count();
8.5.3 跳跃窗口(Hopping)
TimeWindows hoppingWindow = TimeWindows.ofSizeAndGrace(
Duration.ofMinutes(1), // 窗口大小
Duration.ofSeconds(30) // 滑动间隔
);
8.5.4 会话窗口(Session)
import org.apache.kafka.streams.kstream.SessionWindows;
SessionWindows sessionWindow = SessionWindows.with(
Duration.ofMinutes(5) // 静默期
);
KTable<Windowed<String>, Long> sessionCounts = textLines
.groupBy((key, value) -> value)
.windowedBy(sessionWindow)
.count();
8.5.5 窗口抑制(Suppression)
import org.apache.kafka.streams.kstream.Suppressed;
wordCounts
.suppress(Suppressed.untilWindowCloses(
BufferConfig.unbounded()
))
.toStream()
.to("suppressed-output");
8.6 交互查询
8.6.1 概述
交互查询允许直接从应用程序的状态存储中查询数据,无需通过 Kafka 主题。
8.6.2 查询状态存储
import org.apache.kafka.streams.StoreQueryParameters;
import org.apache.kafka.streams.state.KeyValueStore;
KafkaStreams streams = ...;
// 查询
KeyValueStore<String, Long> store =
streams.store(
StoreQueryParameters.fromNameWithType(
"word-count-store",
QueryableStoreType.keyValueStore()
)
);
Long count = store.get("kafka");
8.6.3 远程查询
// 分布式查询
ReadOnlyKeyValueStore<String, Long> store =
streams.store(
StoreQueryParameters.fromNameWithType(
"word-count-store",
QueryableStoreType.keyValueStore()
).withStaleStores()
);
8.7 配置与调优
8.7.1 应用配置
# 基础配置
application.id=my-streams-app
bootstrap.servers=localhost:9092
# 状态目录
state.dir=/tmp/kafka-streams
# 处理保证
processing.guarantee=exactly_once_v2
# 提交间隔
commit.interval.ms=1000
8.7.2 并行度配置
# Streams 线程数(建议与 CPU 核心数相同)
num.stream.threads=4
# 分区缓冲区
buffered.records.per.partition=1000
# 任务超时
task.timeout.ms=300000
8.7.3 内存配置
# 缓存大小
cache.max.bytes.buffering=10485760 # 10MB
# RocksDB 配置
rocksdb.config.setter=com.example.RocksDBConfig
8.7.4 生产者配置
# 生产者配置(继承自 Kafka Producer)
producer.max.block.ms=5000
producer.retries=3
8.7.5 消费者配置
# 消费者配置(继承自 Kafka Consumer)
consumer.max.poll.records=500
consumer.session.timeout.ms=30000
8.8 测试
8.8.1 TopologyTestDriver
import org.apache.kafka.streams.TopologyTestDriver;
import org.apache.kafka.streams.TestInputTopic;
import org.apache.kafka.streams.TestOutputTopic;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.common.serialization.Serdes;
public class WordCountTest {
public void testWordCount() {
// 构建拓扑
StreamsBuilder builder = new StreamsBuilder();
// ... 添加处理逻辑
Topology topology = builder.build();
// 创建 TestDriver
try (TopologyTestDriver driver =
new TopologyTestDriver(topology, getConfig())) {
// 创建测试主题
TestInputTopic<String, String> input =
driver.createInputTopic(
"input",
Serdes.String(),
Serdes.String()
);
TestOutputTopic<String, Long> output =
driver.createOutputTopic(
"output",
Serdes.String(),
Serdes.Long()
);
// 发送输入
input.pipeInput("hello kafka");
// 验证输出
assertEquals(Long.valueOf(1), output.readValue());
// 发送更多输入
input.pipeInput("hello streams");
input.pipeInput("kafka streams");
// 验证输出
assertEquals(Long.valueOf(2), output.readValue());
}
}
}
8.8.2 使用 TestUtils
import org.apache.kafka.streams.TestUtils;
long timeout = 60000; // 60秒
// 等待指定记录
TestUtils.waitForCondition(
() -> output.getQueueSize() > 0,
timeout,
"等待输出"
);
8.9 部署运维
8.9.1 应用打包
Maven:
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>3.5.1</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<configuration>
<transformers>
<transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"/>
<transformer implementation="org.apache.kafka.streams.health.HintTransformer"/>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
8.9.2 运行应用
# 直接运行
java -jar my-streams-app.jar
# Kubernetes 部署
# 需配置合理的资源限制
# liveness/readiness probe
8.9.3 健康检查
// 实现 Health Indicator
public class StreamsHealthIndicator implements HealthIndicator {
private final KafkaStreams streams;
@Override
public Health health() {
KafkaStreams.State state = streams.state();
if (state == KafkaStreams.State.RUNNING) {
return Health.up().build();
} else if (state == KafkaStreams.State.REBALANCING) {
return Health.unknown().build();
} else {
return Health.down().build();
}
}
}
8.9.4 主题管理
# Kafka Streams 创建的内部主题
# - <app-id>-KSTREAM-AGGREGATE-STATE-STORE-0000000001-changelog
# - <app-id>-KSTREAM-MAP-0000000002-repartition
# 这些主题在应用启动时自动创建
# 可通过 Kafka Streams 配置自定义前缀
8.9.5 重置应用
# 使用 reset-tool 重置偏移量
bin/kafka-streams-application-reset.sh \
--application-id my-streams-app \
--bootstrap-servers localhost:9092 \
--input-topics input-topic \
--execute
附录:常见问题
Q1: 如何处理乱序?
// 使用窗口和抑制
KTable<String, Long> counts = source
.groupBy((key, value) -> value)
.windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(1), Duration.ofSeconds(30)))
.count()
.suppress(Suppressed.untilWindowCloses(BufferConfig.unbounded()));
Q2: 如何处理背压?
# 增加缓冲区
cache.max.bytes.buffering=10485760
# 调整拉取大小
consumer.max.poll.records=500
Q3: 如何实现 Exactly-Once?
// 配置 Exactly-Once V2
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
StreamsConfig.EXACTLY_ONCE_V2);