架构师的登山之路|第七站:Kafka + Flink 如何打造一个实时计算链路?
架构师的登山之路|第七站:Kafka + Flink 如何打造一个实时计算链路?
上一站我们聊完了从 Hadoop 到 Flink 的大数据组件学习路线,这一站我们正式走进实时计算的「实战工地」:用 Kafka + Flink 打造一条能落地、好维护、可扩展的实时计算链路。
一、为什么实时计算离不开「Kafka + Flink」组合?
大数据体系里,实时计算负责回答这样几类问题:
- 某个指标现在是多少?
- 某个用户刚刚做了什么?
- 某条异常行为需不需要马上告警或拦截?
- 系统有没有出现异常波动,需不需要立刻处理?
要做好这些事,大致需要两样东西:
- 一条稳定、高吞吐、可扩展的实时数据管道 —— Kafka;
- 一个低延迟、有状态、支持复杂时间语义的实时计算引擎 —— Flink。
简单一句话:Kafka 管“运”,Flink 管“算”。先搞清楚各自职责,再谈如何配合。
二、Kafka:实时数据世界的「高速公路」
可以把 Kafka 想象成一家公司的「总线 + 物流系统」:任何系统要产生日志、埋点、事件,都往这条总线里丢;任何需要消费这些数据的系统,都从总线里订阅。
2.1 核心概念一图掌握
- Broker:Kafka 服务器节点,每个 Broker 存一部分数据。
- Topic:逻辑上的「主题」,按业务维度划分,例如
user-behavior,order-log。 - Partition:Topic 的分片,真正承载数据和并行度的单位。
- Producer:消息生产者,例如埋点 SDK、日志采集器。
- Consumer / Consumer Group:消息消费者和消费组,同一组内的消费者共享一个 Topic 的消费进度。
- Offset:每条消息在分区中的位置,用来记录消费进度。
记住一句话:吞吐量靠 Partition,可用性靠副本(Replication),消费水平扩展靠 Consumer Group。
2.2 设计 Kafka 时要考虑什么?
-
分区数(Partitions)
- 分区越多,并行度越高,但管理成本也越大;
- 一般建议:按预期最大并行度 + 增长空间来估算,而不是拍脑袋。
-
副本因子(Replication Factor)
- 通常设置为 2 或 3;
- 越高越可靠,但占用磁盘更多。
-
消息保留策略
- 按时间:如保留 7 天;
- 按空间:如保留到 Topic 达到多少 GB;
- 具体根据业务恢复需求、合规要求决定。
-
生产与消费策略
- 生产端:
acks、重试、压缩(lz4/snappy); - 消费端:自动提交 vs 手动提交 offset、幂等处理。
- 生产端:
2.3 Kafka 常见使用场景
- 行为埋点、日志收集;
- 交易流水、监控指标流;
- IoT 设备数据上报;
- CDC(变更数据捕获),把数据库变更实时投递出去。
三、Flink:实时计算链路的「大脑」
有了 Kafka 这条高速公路,还需要一个能在数据流动过程中「实时思考」的引擎,Flink 就是干这个的。
3.1 为什么是 Flink?
- 毫秒级延迟:真正意义上的流式实时处理;
- 事件时间(Event Time)支持:可以处理乱序、延迟到达的数据;
- 强大的状态管理:可以维护大规模有状态计算(如实时窗口聚合、实时画像);
- 一致性保障:通过 Checkpoint + 两阶段提交实现 Exactly-once。
3.2 必须掌握的几个核心概念
-
时间语义
- Processing Time:算子执行该条数据时的机器时间;
- Event Time:事件真实发生的时间(需要数据里有时间字段);
- Watermark(水位线):用来标记「某个时间点之前的事件基本都到了」。
-
窗口(Window)类型
- 滚动窗口(Tumbling Window):如每 5 分钟滚动统计;
- 滑动窗口(Sliding Window):如每 1 分钟滑动统计过去 5 分钟;
- 会话窗口(Session Window):基于「空闲间隔」划分会话,如用户会话行为。
-
状态(State)与容错
- 键控状态(Keyed State)、算子状态(Operator State);
- Checkpoint:周期性保存某个时刻所有算子的状态快照;
- Savepoint:人为触发,常用于升级、迁移、回滚。
3.3 Kafka + Flink 的典型集成方式
- Source:Flink 通过 Kafka Connector 消费 Topic;
- 算子链(Operators):过滤、清洗、聚合、连接、富函数、自定义逻辑等;
- Sink:把结果写入下游:
- HBase / Redis:实时查询;
- Elasticsearch / ClickHouse:实时检索与分析;
- Kafka:继续派发到下游系统。
四、一条典型实时计算链路长什么样?
以「实时用户行为看板」为例,典型链路可以这样设计:
-
数据产生
- 前端 / App 埋点,上报「PV、UV、点击事件、曝光事件」;
- Nginx / 应用服务记录访问日志。
-
数据采集与进入 Kafka
- 通过埋点 SDK、Fluentd、Filebeat 等采集日志;
- 写入 Kafka 的
user-behaviorTopic。
-
Flink 实时处理
- Source:消费
user-behavior; - 清洗:过滤脏数据、补充字段、格式统一;
- 维度关联:如补充用户归属、渠道信息(可以通过维表 join);
- 窗口聚合:
- 每分钟实时统计 PV/UV;
- 每 5 分钟统计各渠道转化率;
- 对异常波动设置告警规则。
- Source:消费
-
结果写入与展示
- 写入 Redis:实时大屏接口直接从 Redis 读;
- 写入 ClickHouse / ES:支持灵活维度分析;
- 写入 Kafka 其他 Topic:供风控 / 推荐等其他系统继续消费。
五、设计 Kafka + Flink 链路的关键决策点
5.1 Topic 与分区规划
- 尽量按业务维度 + 环境划分:
- 例:
user-behavior-prod、order-log-prod;
- 例:
- 分区数影响:
- 吞吐与并行度;
- 数据倾斜的可能性(例如按 userId 取模时要注意热点用户)。
5.2 Flink 并行度与资源规划
- 一般建议:Flink 任务并行度 ≈ Kafka 分区数 或者其整数倍;
- 对于极不均匀的 Key(如大 V 用户),考虑:
- 二级 key 分组;
- 或者使用类似「打散 + 再聚合」的模式。
5.3 状态与 Checkpoint 策略
- 状态较大的作业:
- 建议使用 RocksDB State Backend;
- 合理设置 Checkpoint 间隔(如 1 分钟 / 5 分钟);
- Checkpoint 保存位置:
- HDFS / 对象存储(如 S3 兼容存储);
- 与外部系统互操作时:
- 优先使用支持两阶段提交的 Sink;
- 否则通过「幂等写入 + 事务标记」来保障结果可接受。
5.4 监控与告警
- 指标维度:
- Kafka Topic 堆积情况;
- Flink Job 的延迟、反压(Backpressure)、Checkpoint 成功率;
- 下游存储(Redis/HBase/ES)响应时间与错误率。
- 工具:
- Prometheus + Grafana;
- Flink 自带 Web UI、Kafka Manager 等。
六、实战里常见的几个坑
6.1 Kafka 分区太少 or 不可扩展
- 一开始只给了 1~2 个分区,后面流量上来并行度不够;
- 虽然 Kafka 支持增加分区,但现有数据的顺序语义会被打乱。
建议:按「峰值预估 + 一定裕量」来规划分区,不要只看当前流量。
6.2 Flink 状态失控,Checkpoint 卡死
- 状态越积越大,Checkpoint 时间越来越长;
- Checkpoint 太频繁导致 IO 压力过大,作业出现反压。
建议:
- 控制窗口大小、避免长时间未清理的状态;
- 对极长生命周期的状态,考虑分层:一部分在实时,一部分离线归档。
6.3 只考虑计算,不考虑下游能不能扛
- Flink 算得飞快,结果写入下游(HBase/ES/Redis)时被打满;
- 导致延迟急剧上升,甚至作业反复重启。
建议:
- 为 Sink 侧增加限流、重试、熔断机制;
- 对热点 key 做打散;
- 必要时引入中间缓冲层(如写回 Kafka / Pulsar 再异步入库)。
七、一份可以直接套用的 Kafka + Flink 规划清单
在你准备上一个实时计算项目时,可以按这份清单过一遍:
-
业务目标是否真的需要实时?
- 需要秒级 / 分钟级;
- T+1 足够?能否先用离线方案验证价值?
-
数据从哪里来?落到哪些 Topic?
- 是否已经有统一的埋点 / 日志体系?
- Topic 命名与分区、保留策略是否规划好?
-
实时作业要算什么?输出给谁?
- 指标、维度、窗口类型;
- 下游系统是看板、风控、推荐还是告警?
-
Flink Job 如何拆分?
- 是否拆成多个独立作业,降低耦合;
- 是否有公共逻辑需要抽成库 / 模板?
-
状态、Checkpoint 与容错策略?
- 状态是否会持续增长?如何清理?
- Checkpoint 周期、超时时间、存储位置?
-
监控、告警与排错路径?
- 需要哪些关键指标;
- 事故发生后,谁负责排查、该看哪些图表?
八、总结:Kafka + Flink 是实时计算的「黄金组合」
这一站我们从架构师视角,把 Kafka + Flink 组合在实时链路中的角色和设计要点串了起来:
- Kafka 提供 高吞吐、可扩展、持久化的实时数据管道;
- Flink 提供 低延迟、有状态、支持复杂时间语义的实时计算能力;
- 两者组合起来,能支撑从实时看板、实时风控,到实时推荐、实时监控的一整套业务场景。
关键不在于「我会多少 API」,而在于:
能不能设计出一条稳定、可演进、遇到问题好排查的实时链路。
下一站预告
实时计算链路搭好之后,数据只是「跑起来」了,还谈不上「用得好」。
下一站,我们会把视角从「链路」切换到「治理与协同」:
聊聊 数据治理、元数据管理、数据中台 这些经常被提起却容易被误解的概念,
以及它们与一线开发、架构工作之间到底是什么关系。
敬请期待。
更多推荐

所有评论(0)