📌 本文是「Flink 教程」系列第 七 篇 · 系列目录
A. 配置参数参考
A.1 核心运行时参数
| 参数 | 默认值 | 是否必填 | 说明 |
|---|---|---|---|
parallelism.default | 1 | ❌ | 默认并行度 |
taskmanager.numberOfTaskSlots | 1 | ❌ | 每个 TaskManager 的 Slot 数 |
taskmanager.memory.process.size | - | ✅ | TaskManager 总内存 |
jobmanager.memory.process.size | - | ✅ | JobManager 总内存 |
A.2 Checkpoint 参数
| 参数 | 默认值 | 是否必填 | 说明 |
|---|---|---|---|
execution.checkpointing.interval | - | ❌ | Checkpoint 间隔(毫秒) |
execution.checkpointing.mode | EXACTLY_ONCE | ❌ | Checkpoint 语义 |
execution.checkpointing.timeout | 10min | ❌ | Checkpoint 超时时间 |
state.backend | - | ✅ | 状态后端类型 |
state.checkpoints.dir | - | ✅ | Checkpoint 存储路径 |
A.3 RocksDB 参数(可选)
| 参数 | 默认值 | 说明 |
|---|---|---|
state.backend.rocksdb.memory.managed | true | 是否启用托管内存 |
state.backend.rocksdb.localdir | - | RocksDB 本地存储路径(推荐 SSD) |
state.backend.rocksdb.checkpoint.transfer.thread.num | 1 | Checkpoint 上传线程数 |
B. 依赖版本
- Flink: 1.14.4
- Scala: 2.12 (对应flink-streaming-java_2.12)
- Java: 8+
C. 扩展阅读
D. 生产环境常见问题 FAQ(补充)
本章节基于《Flink原理、实战与性能优化》与实际生产经验整理
D.1 ⏰ Watermark 相关问题
Q: Watermark 不前进怎么办?
- 检查时间戳分配器是否正确
- 确保事件时间戳 > 0
- 检查无序边界设置是否合理
Q: 窗口不触发怎么办?
- 确保已分配时间戳和 Watermark
- 检查 Watermark 是否超过窗口结束时间
- 检查是否使用了正确的时间语义
D.2 💾 Checkpoint 相关问题
Q: Checkpoint 总是失败?
- 检查 Checkpoint 超时时间是否足够
- 检查状态大小是否超出限制
- 检查存储系统(HDFS/S3)是否正常
- 考虑启用增量 Checkpoint
Q: Checkpoint 时间过长?
- 考虑使用 RocksDBStateBackend
- 启用非对齐 Checkpoint
- 优化状态后端配置
- 监控 Checkpoint Delay Time
D.3 🗄️ RocksDB 相关问题
Q: RocksDB 写入性能差?
- 增加写缓冲区数量和大小
- 调整压缩策略
- 使用托管内存
- 确保使用 SSD 存储
Q: 状态数据超过 2GB?
- 拆分大状态为多个 Keyed State
- 使用状态 TTL 清理过期数据
- 考虑业务逻辑优化
D.4 🎯 Exactly-Once 相关问题
Q: 如何实现端到端 Exactly-Once?
- Source 端:从 Checkpoint 恢复消费位置
- Flink 端:EXACTLY_ONCE Checkpoint + 两阶段提交
- Sink 端:支持事务的连接器(如 Kafka 事务)
Q: Kafka Sink 重复数据?
- 检查 FlinkKafkaProducer 的 Semantic 配置
- 确保使用 EXACTLY_ONCE 模式
- 检查事务超时配置
文档更新时间:2026年 基于 Apache Flink 1.14 官方文档 + 《Flink原理、实战与性能优化》补充
标签: #flink #附录 #参考 #FAQ