这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Kafka 的强大一半来自生态。本章介绍三个官方系工具:Connect 解决“数据进出 Kafka”,Streams 解决“流处理”,ksqlDB 让你用 SQL 写流式应用。
23.1 Kafka Connect
是什么
Connect 是一个数据集成框架,以插件方式运行 connector,无需写代码即可在 Kafka 与外部系统之间搬运数据:
Source Connector: 外部系统 -> Kafka
JDBC / Debezium(CDC) / File / S3 ...
Sink Connector: Kafka -> 外部系统
JDBC / Elasticsearch / S3 / ClickHouse ...
核心概念
| 概念 | 说明 |
|---|---|
| Connector | 逻辑任务定义(连接哪个库、读哪些表) |
| Task | 实际并行工作的执行单元 |
| Worker | 运行 connector/task 的进程 |
| Converter | 消息格式转换(JSON/Avro/Protobuf) |
| Offset | source 端进度(如表的 binlog 位点),存于内部 topic |
| REST API | 管理 connector 的 HTTP 接口 |
两种模式
standalone:单进程,适合开发测试
distributed:多 worker,任务自动均衡与容错,生产必选
启动(distributed)
bin/connect-distributed.sh config/connect-distributed.properties
关键配置:
bootstrap.servers=localhost:9092
group.id=connect-cluster
offset.storage.topic=connect-offsets
config.storage.topic=connect-configs
status.storage.topic=connect-status
offset.storage.replication.factor=3
config.storage.replication.factor=3
status.storage.replication.factor=3
key.converter=org.apache.kafka.connect.json.JsonConverter
value.converter=org.apache.kafka.connect.json.JsonConverter
用 REST 创建 JDBC Source
curl -X POST http://localhost:8083/connectors \
-H 'Content-Type: application/json' -d '{
"name": "jdbc-orders-source",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
"tasks.max": "4",
"connection.url": "jdbc:mysql://db:3306/shop",
"connection.user": "reader",
"connection.password": "***",
"table.whitelist": "orders",
"mode": "incrementing",
"incrementing.column.name": "id",
"topic.prefix": "cdc-"
}
}'
管理命令:
curl localhost:8083/connectors # 列出
curl localhost:8083/connectors/jdbc-orders-source/status
curl -X DELETE localhost:8083/connectors/jdbc-orders-source
选型建议
- CDC 场景优先 Debezium(binlog 精确、对库压力小);
- JDBC source 的
incrementing/timestamp模式对删除与更新历史支持有限; - sink 端要考虑幂等(upsert by key),否则重平衡时可能重复写;
- 生产环境 connector 配置进版本库,REST 变更走 CI/CD。
23.2 Kafka Streams
是什么
Streams 是一个 Java 库(不是独立服务),嵌入应用即可做有状态流计算:
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> lines = builder.stream("input-topic");
KTable<String, Long> counts = lines
.flatMapValues(v -> Arrays.asList(v.toLowerCase().split("\\W+")))
.groupBy((k, word) -> word)
.count(Materialized.as("word-counts"));
counts.toStream().to("output-topic", Produced.with(Serdes.String(), Serdes.Long()));
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();
配置:
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "wordcount-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
application.id 同时充当消费组 ID 与状态存储前缀,同一应用的所有实例必须一致。
KStream vs KTable
| 抽象 | 语义 | 类比 |
|---|---|---|
| KStream | 事件流,每条都重要 | append 日志 |
| KTable | 按 key 聚合的最新状态 | 可更新的表 |
| GlobalKTable | 全实例广播的表 | 小维度表 |
经典组合:订单流(KStream)join 用户表(KTable),实时补全维度。
窗口
clicks.groupByKey()
.windowedBy(TimeWindows.ofSizeAndGrace(Duration.ofMinutes(5), Duration.ofSeconds(10)))
.count();
- Tumbling:固定不重叠;
- Hopping:按步长滑动,可能重叠;
- Session:按活动间隔动态聚合(用户会话);
- Grace period:允许迟到数据,超期进入弃用流。
状态与容错
状态存本地 RocksDB,并后台备份到 changelog topic:
实例宕机 -> 新实例接管分区 -> 从 changelog 重建本地状态 -> 继续处理
因此 Streams 应用会自动创建 <app-id>-<store>-changelog 主题,不要手工删除。
适用场景
- 实时聚合、监控规则、风控特征;
- 流表 join 补维;
- 数据清洗与路由; 需要 exactly-once 的 Kafka-to-Kafka 管道。
复杂拓扑、事件时间语义、大规模状态管理时,再评估 Flink(第 25 章)。
23.3 ksqlDB
ksqlDB 把流处理变成 SQL,适合快速验证与轻量场景。
-- 建流
CREATE STREAM orders (
order_id VARCHAR KEY,
user_id VARCHAR,
amount DECIMAL(12, 2)
) WITH (
KAFKA_TOPIC = 'order-events',
VALUE_FORMAT = 'JSON'
);
-- 每 5 分钟统计各用户消费
CREATE TABLE user_spend AS
SELECT user_id,
COUNT(*) AS order_cnt,
SUM(amount) AS total
FROM orders
WINDOW TUMBLING (SIZE 5 MINUTES)
GROUP BY user_id
EMIT CHANGES;
-- 持续查询(推模式)
SELECT * FROM user_spend
WHERE total > 10000
EMIT CHANGES;
优点:上手快、声明式、内置服务化。注意:复杂逻辑、精细状态控制与大型作业,仍建议 Streams/Flink。
23.4 选型对比
| 需求 | 推荐 |
|---|---|
| 数据库/日志/对象存储进出 Kafka | Connect |
| 库表级 CDC | Debezium + Connect |
| Java 团队的轻中量流处理 | Streams |
| SQL 快速分析/原型 | ksqlDB |
| 重状态、事件时间、大规模作业 | Flink |
| 复杂图计算/批流一体 | Spark/Flink |
本章小结
- Connect 用 connector + worker + 内部 topic 实现高可用数据集成;
- Streams 是库形态的流处理框架,KTable/窗口/changelog 是三大核心;
- ksqlDB 用 SQL 降低流处理门槛,适合轻量与服务化查询;
- 生态选型先看团队栈与运维能力,再谈功能上限。
思考题
- Connect 的 source offset 存在哪里?为什么这比存在内存里可靠?
- Streams 的 changelog topic 删除后会发生什么?如何恢复?
- 用 Streams 实现一个“每 1 分钟统计各接口 5xx 次数并告警”的拓扑,写出伪代码。