跳到正文
SL Blog 技术探索 · 工程实践 · AI 时代思考
返回
📨 Kafka 3.9 教程 · 9 / 10 查看系列简介 →

Kafka 3.9 教程(八):Kafka Streams

文章目录
  1. 目录
  2. 8.1 简介
  3. 8.1.1 什么是 Kafka Streams?
  4. 8.1.2 与其他流处理框架对比
  5. 8.1.3 应用场景
  6. 8.2 核心概念
  7. 8.2.1 KStream
  8. 8.2.2 KTable
  9. 8.2.3 GlobalKTable
  10. 8.2.4 State Store
  11. 8.2.5 Topology
  12. 8.3 DSL API
  13. 8.3.1 基本使用
  14. 8.3.2 无状态操作
  15. 8.3.3 有状态操作
  16. 8.3.4 表操作
  17. 8.3.5 连接操作
  18. 8.3.6 输出操作
  19. 8.3.7 Serde(序列化/反序列化)
  20. 8.3.8 时间语义
  21. 8.3.9 命名主题与资源管理
  22. 8.4 Processor API
  23. 8.4.1 基本用法
  24. 8.4.2 构建拓扑
  25. 8.5 窗口操作
  26. 8.5.1 窗口类型
  27. 8.5.2 滚动窗口(Tumbling)
  28. 8.5.3 跳跃窗口(Hopping)
  29. 8.5.4 会话窗口(Session)
  30. 8.5.5 窗口抑制(Suppression)
  31. 8.6 交互查询
  32. 8.6.1 概述
  33. 8.6.2 查询状态存储
  34. 8.6.3 远程查询
  35. 8.7 配置与调优
  36. 8.7.1 应用配置
  37. 8.7.2 并行度配置
  38. 8.7.3 内存配置
  39. 8.7.4 生产者配置
  40. 8.7.5 消费者配置
  41. 8.8 测试
  42. 8.8.1 TopologyTestDriver
  43. 8.8.2 使用 TestUtils
  44. 8.9 部署运维
  45. 8.9.1 应用打包
  46. 8.9.2 运行应用
  47. 8.9.3 健康检查
  48. 8.9.4 主题管理
  49. 8.9.5 重置应用
  50. 附录:常见问题
  51. Q1: 如何处理乱序?
  52. Q2: 如何处理背压?
  53. Q3: 如何实现 Exactly-Once?
  54. 下一步学习

📌 本文是「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 StreamsApache FlinkSpark 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);
}

状态存储类型

存储类型接口说明
KeyValueStoreKeyValueStore键值存储
WindowStoreWindowStore窗口键值存储
SessionStoreSessionStore会话存储

状态容错

# 启用 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转换键值对KStreamKStream
mapValues只转换值KStream/KTable相同
flatMap展开转换KStreamKStream
flatMapValues展开值KStreamKStream
filter过滤KStream/KTable相同
selectKey重新指定 KeyKStreamKStream
merge合并流KStreamKStream
// 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按键分组KStreamKGroupedStream
groupByKey按现有键分组KStreamKGroupedStream
// 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);

下一步学习


RAG 智能问答

针对本文继续提问:《Kafka 3.9 教程(八):Kafka Streams》