🌊 第 17 章 · Kafka Streams 与 ksqlDB 可视化
九个模块:① 流-表二元性 ② 4 种窗口对比 ③ KStream-KTable Join ④ WordCount Topology
⑤ KStream/KTable/GlobalKTable 选型 ⑥ 三种 Join 对比 ⑦ 时间语义 ⑧ State Store 容灾 ⑨ EOS v2 ⑩ ksqlDB 速查 ⑪ 常见坑。
① 流-表二元性(Stream-Table Duality)
左:KStream(事件流,每条都是一个事实);右:KTable(按 key 取最新值的当前状态)。点击「+ 加事件」感受两者关系。
KTable(balances)= sum by key
② 四种窗口类型对比动画
点击上方按钮选择窗口类型。
③ KStream-KTable Join(流加维度表)
KTable 维护用户档案(按 user_id 取最新);KStream 进来订单时用 user_id 查 KTable 加上用户名 / 国家。
④ WordCount Topology 可视化
输入一句话,观察消息在拓扑节点之间流动。
📨
Source
learn.text
→
🔪
flatMapValues
split words
→
🔑
groupBy
by word
→
🔢
count
State Store
→
📤
Sink
learn.wc-output
State Store: counts-store
Changelog Topic(compacted)
(等待输入…)
⑤ KStream / KTable / GlobalKTable 选型
| 抽象 | 语义 | 分区 | 典型用途 | 反例 |
| KStream |
事件流,每条都是独立事实 |
按消息 key |
点击日志、订单事件、IoT 上报 |
不要存最新状态(每条都保留) |
| KTable |
按 key 取最新值的实时快照 |
按 key(与 KStream 同分区策略) |
用户档案、商品价格、配置项 |
不要做事件计数(应改 KStream → count) |
| GlobalKTable |
全集群每个 task 都拥有完整副本 |
不分区(整张表广播) |
小型维度表、字典(≤ 100 MB) |
不要给百万行大表用,每个实例都要装一份 |
💡 KStream-KTable Join 要求两边同 key 同分区;KStream-GlobalKTable Join 不需要 co-partitioning,可按任意外键查找。
⑥ 三种 Join 对比
| Join 类型 | 是否要 window | 触发 | 典型场景 |
| KStream-KStream |
✅ 必须(如 5min) |
两边都来才匹配 |
下单 + 支付,5 分钟内匹配;点击 + 转化 |
| KStream-KTable |
❌ 不需要 |
左边来一条就查右表 |
订单流 + 用户档案表 → 订单加用户名 |
| KTable-KTable |
❌ 不需要 |
任何一边变化都触发 |
账户余额 + 利率表 → 实时利息 |
| KStream-GlobalKTable |
❌ 不需要 |
左边来一条查全局表(可用外键) |
事件流 + 全量小型字典(国家码、币种) |
// KStream-KTable join 示例(不需要 window)
KStream<String, Order> orders = builder.stream("orders");
KTable<String, User> users = builder.table("users");
orders.join(users, (order, user) -> enrich(order, user))
.to("orders.enriched");
⑦ 三种时间语义
| 类型 | 来源 | 优点 | 缺点 |
| Event Time |
消息体里的业务时间(如 ts 字段) |
正确反映业务顺序,乱序也可处理 |
需要自定义 TimestampExtractor,要处理 watermark / late event |
| Ingestion Time |
Broker 接收消息时的 server timestamp |
无需用户配置;接近事件时间 |
客户端积压时会失真 |
| Processing Time |
Streams 处理消息那一刻的本地系统时间 |
实现简单,无依赖 |
重放历史数据时窗口完全错位 |
// 自定义 Event Time 抽取
public class JsonTsExtractor implements TimestampExtractor {
public long extract(ConsumerRecord<Object,Object> r, long partitionTime) {
long ts = parseJson((String) r.value()).getLong("event_ts");
return ts > 0 ? ts : partitionTime; // 兜底:用上一次的时间
}
}
props.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG, JsonTsExtractor.class);
💡 配置 grace.ms 接受迟到事件,但越长状态越大;金融场景常 grace=5min,普通分析 grace=10s。
⑧ State Store + Changelog Topic(容灾)
RocksDB 在本地磁盘存状态;每次 put 同步写入 changelog topic(compacted)。Task 迁移到新 Worker 时,从 changelog 重放重建本地 store。
[Task-0] -- put(k=alice, v=10) --> RocksDB (本地)
|
+--> producer.send(changelog-topic, k=alice, v=10)
崩溃后新 Worker 接手 Task-0:
1. 创建空 RocksDB
2. 从 changelog topic 头到尾 replay → 状态完全恢复
3. 切换为正常处理消息
⚠ 优化:
- changelog 是 compacted topic(log.cleanup.policy=compact),同 key 只保留最新
- num.standby.replicas=1:另一台 Worker 实时跟随 changelog → 故障接管秒级
- 关闭:StoreBuilder.withLoggingDisabled() → 只本地、不容灾(仅适合纯衍生数据)
⑨ EOS v2 在 Streams 中的开启
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
StreamsConfig.EXACTLY_ONCE_V2); // ← 一行打开
props.put(StreamsConfig.REPLICATION_FACTOR_CONFIG, 3); // 内部 topic ≥ 3
| 版本 | 差别 |
| EOS v1(已废弃) | 每个 input partition 单独事务,开销大 |
| EOS v2(推荐) | 每个 thread 共享一个 producer + 事务,性能接近 ALOS(At Least Once Semantic) |
⚠ 下游消费者必须 isolation.level=read_committed,否则会看到 abort 的中间值。
⑩ ksqlDB 常见 SQL 速查
-- 创建 STREAM(基于 Topic)
CREATE STREAM clicks (user_id VARCHAR, page VARCHAR, ts BIGINT)
WITH (KAFKA_TOPIC='learn.clicks', VALUE_FORMAT='JSON', TIMESTAMP='ts');
-- 创建 TABLE(按 key 取最新)
CREATE TABLE users (user_id VARCHAR PRIMARY KEY, name VARCHAR, country VARCHAR)
WITH (KAFKA_TOPIC='learn.users', VALUE_FORMAT='JSON');
-- 滚动窗口聚合:每分钟每用户 PV
CREATE TABLE pv_per_min AS
SELECT user_id, COUNT(*) AS pv
FROM clicks
WINDOW TUMBLING (SIZE 1 MINUTE)
GROUP BY user_id
EMIT CHANGES;
-- Stream-Table Join(流加维度)
CREATE STREAM enriched AS
SELECT c.user_id, u.name, u.country, c.page
FROM clicks c LEFT JOIN users u ON c.user_id = u.user_id
EMIT CHANGES;
-- Pull Query(点查 TABLE 当前值)
SELECT * FROM users WHERE user_id = 'u-1';
-- 持续查询管理
SHOW QUERIES;
TERMINATE CSAS_PV_PER_MIN_3;
DROP TABLE pv_per_min DELETE TOPIC;
⑪ 常见坑速查
| 坑 | 表现 | 修复 |
| Co-partitioning 不一致 |
KStream-KTable join 出现 NPE / 漏匹配 |
两边 partition 数必须相同;上游不一致就先 repartition() |
| 用 KStream.count() 直接当库存 |
每次重启数都翻倍 |
应该用 KTable + 业务事件 = 增量更新;或定期快照 |
| State Store 巨大但没启 standby |
Worker 故障切换要恢复几十分钟 |
num.standby.replicas=1 + 监控 changelog 大小 |
| Window 不设 grace |
迟到 1 秒的事件全丢 |
.grace(Duration.ofSeconds(10)) |
| Application.id 改了 |
所有 changelog/repartition topic 重建,状态全失 |
application.id = 业务身份证,永远不要随意改 |
| EOS v2 没改下游 |
下游看到 abort 的脏数据 |
所有下游 Consumer / ksqlDB / Connect Sink 全部 read_committed |
⑫ Streams / Flink / Spark Streaming 选型
| 维度 | Kafka Streams | Flink | Spark Structured Streaming |
| 部署形态 | 普通 Java 应用(Jar) | 独立集群(JM/TM) | Spark 集群(Driver/Executor) |
| 外部数据源 | 仅 Kafka | 多种 source/sink | 多种 |
| 状态后端 | RocksDB + changelog | RocksDB / 内存 + checkpoint | HDFS checkpoint |
| EOS | v2,开箱即用 | checkpoint + 2PC sink | v3 后支持 |
| 窗口 | Tumbling / Hopping / Session / Sliding | 同 + 自定义触发器,更灵活 | Tumbling / Sliding / Session |
| SQL | ksqlDB(独立服务) | Flink SQL / Table API | Spark SQL |
| 定位 | Kafka 内嵌轻量级流处理 | 纯流计算引擎,金融实时首选 | 批流统一,离线为主 |
💡 选型口诀:
· 处理 100% 在 Kafka 内部、想简单部署 → Kafka Streams / ksqlDB
· 多源 + 复杂窗口 + 大状态 + 金融级 EOS → Flink
· 已有 Spark 离线团队、需要批流统一 → Spark Structured Streaming
⑬ Streams 关键调优参数
| 参数 | 默认 | 建议 | 原因 |
num.stream.threads | 1 | = max(input partition 数) | 每线程跑若干 task,增大并行度 |
num.standby.replicas | 0 | 1 | 状态可秒级切换 |
cache.max.bytes.buffering | 10MB | 50~200MB | 聚合 dedupe,减少下游 RPS |
commit.interval.ms | 30s(EOS 100ms) | EOS 200~500ms | 降低 EOS 模式吞吐损耗 |
processing.guarantee | at_least_once | exactly_once_v2(金融) | EOS v2 性能损耗 < 10% |
buffered.records.per.partition | 1000 | 10000+ | 抗瞬时峰值 |
state.dir | /tmp/kafka-streams | 独立 SSD 目录 | RocksDB 强 IO |
📚 第 17 章配套:17_streams.md · code/wordcount_streams.java · code/wordcount_faust.py
code/realtime_topn.py · code/windowed_aggregate.py · code/ksql_examples.sql