Flink 与 Kafka 集成
概述
Kafka + Flink 是实时计算的事实组合:Flink 用 Kafka 作数据管道与结果出口。本文讲清 Kafka Source/Sink 连接器、offset 提交与 checkpoint 的关系、动态分区发现,以及 Exactly-Once 写入的端到端一致性方案。
一、Kafka 连接器
1.1 DataStream 方式
java
// Source
DataStream<String> stream = env.addSource(
new FlinkKafkaConsumer<>(
"orders",
new SimpleStringSchema(),
kafkaProps));
// Sink
stream.addSink(
new FlinkKafkaProducer<>(
"result",
new SimpleStringSchema(),
kafkaProps));1.2 Table 方式
sql
CREATE TABLE orders (
order_id BIGINT,
user_id BIGINT,
amount DECIMAL(10, 2),
ts TIMESTAMP(3)
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'kafka:9092',
'properties.group.id' = 'flink-group',
'format' = 'json',
'scan.startup.mode' = 'latest-offset'
);二、Kafka Source
2.1 消费语义
| 参数 | 说明 |
|---|---|
scan.startup.mode | earliest/latest/group-offsets/specific-offsets |
properties.group.id | 消费组 |
properties.bootstrap.servers | 集群地址 |
2.2 启动模式
| 模式 | 说明 |
|---|---|
| earliest | 从最早开始 |
| latest | 从最新开始 |
| group-offsets | 从已提交 offset 恢复 |
| specific-offsets | 指定 offset |
| timestamp | 从时间点开始 |
2.3 分区发现
Flink 定期探测 Topic 新增分区:
| 参数 | 默认 | 说明 |
|---|---|---|
flink.partition-discovery.interval-millis | 关闭 | 分区发现间隔 |
开启后:自动消费新分区(扩容无需重启)三、offset 提交机制
3.1 与 Checkpoint 的关系
Checkpoint 时:
1. Source 保存各分区 offset(状态)
2. 全部算子快照完成
3. 恢复时从快照 offset 重放
auto.commit 与 checkpoint:
checkpoint 开启时自动提交(disable auto)
未开启时才依赖自动提交| 配置 | 说明 |
|---|---|
checkpointing.enabled | 开启则用 checkpoint 管理 offset |
enable.auto.commit | 建议 false |
auto.commit.interval.ms | 自动提交间隔 |
3.2 语义
| 模式 | offset 管理 | 语义 |
|---|---|---|
| Checkpoint 开启 | 快照提交 | 精确一次(配合幂等 Sink) |
| 自动提交 | Kafka 提交 | 至少一次 |
四、Kafka Sink
4.1 写入配置
sql
CREATE TABLE result (...) WITH (
'connector' = 'kafka',
'topic' = 'result',
'format' = 'json',
'sink.partitioner' = 'default',
'sink.semantic' = 'exactly-once'
);| 参数 | 说明 |
|---|---|
sink.partitioner | default/fixed/round-robin/custom |
sink.semantic | at-least-once/exactly-once |
sink.transaction-timeout | 事务超时 |
4.2 分区器
| 分区器 | 行为 |
|---|---|
| default | 按 key 哈希 |
| fixed | 固定分区 |
| round-robin | 轮询 |
| custom | 自定义 |
五、Exactly-Once 写入
5.1 两阶段提交
Flink 的 Kafka Sink 用两阶段提交(2PC):
1. 预提交:Checkpoint 时提交事务
2. 确认:Checkpoint 完成,提交事务
3. 失败:回滚事务5.2 条件
| 条件 | 说明 |
|---|---|
| 开启 Checkpoint | 必须 |
sink.semantic=exactly-once | Sink 事务模式 |
| Kafka 支持事务 | 事务超时配置 |
5.3 端到端一致性
Kafka Source(offset 快照)
→ Flink 状态(checkpoint 恢复)
→ Kafka Sink(两阶段提交)
= 端到端 Exactly-Once| 组件 | 保证 |
|---|---|
| Source | 快照 offset,重放 |
| 引擎 | 状态恢复 |
| Sink | 事务提交 |
六、常见问题
6.1 事务超时
| 问题 | 说明 |
|---|---|
| TransactionTimeout | 默认 1h,需大于 checkpoint 间隔 |
| 冲突 | 与 Kafka broker 事务超时匹配 |
6.2 消费不均衡
| 原因 | 处理 |
|---|---|
| 分区数 < 并行度 | 增加分区 |
| 数据分布不均 | 检查 key 设计 |
| 分区发现未开 | 开启自动发现 |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 重复消费 | checkpoint 与 Sink 语义不匹配,统一精确一次 |
| 新分区不消费 | 开启 partition-discovery |
| 事务失败 | 事务超时与 checkpoint 间隔不匹配 |
| 启动模式错误 | 按需求设置 scan.startup.mode |
| 背压到 Kafka | 下游处理慢,优化或扩容 |