架构师的登山之路|第七站:Kafka + Flink 如何打造一个实时计算链路?

上一站我们聊完了从 Hadoop 到 Flink 的大数据组件学习路线,这一站我们正式走进实时计算的「实战工地」:用 Kafka + Flink 打造一条能落地、好维护、可扩展的实时计算链路。


一、为什么实时计算离不开「Kafka + Flink」组合?

大数据体系里,实时计算负责回答这样几类问题:

  • 某个指标现在是多少?
  • 某个用户刚刚做了什么?
  • 某条异常行为需不需要马上告警或拦截?
  • 系统有没有出现异常波动,需不需要立刻处理?

要做好这些事,大致需要两样东西:

  1. 一条稳定、高吞吐、可扩展的实时数据管道 —— Kafka;
  2. 一个低延迟、有状态、支持复杂时间语义的实时计算引擎 —— 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 时要考虑什么?

  1. 分区数(Partitions)

    • 分区越多,并行度越高,但管理成本也越大;
    • 一般建议:按预期最大并行度 + 增长空间来估算,而不是拍脑袋。
  2. 副本因子(Replication Factor)

    • 通常设置为 2 或 3;
    • 越高越可靠,但占用磁盘更多。
  3. 消息保留策略

    • 按时间:如保留 7 天;
    • 按空间:如保留到 Topic 达到多少 GB;
    • 具体根据业务恢复需求、合规要求决定。
  4. 生产与消费策略

    • 生产端:acks、重试、压缩(lz4/snappy);
    • 消费端:自动提交 vs 手动提交 offset、幂等处理。

2.3 Kafka 常见使用场景

  • 行为埋点、日志收集;
  • 交易流水、监控指标流;
  • IoT 设备数据上报;
  • CDC(变更数据捕获),把数据库变更实时投递出去。

三、Flink:实时计算链路的「大脑」

有了 Kafka 这条高速公路,还需要一个能在数据流动过程中「实时思考」的引擎,Flink 就是干这个的。

3.1 为什么是 Flink?

  • 毫秒级延迟:真正意义上的流式实时处理;
  • 事件时间(Event Time)支持:可以处理乱序、延迟到达的数据;
  • 强大的状态管理:可以维护大规模有状态计算(如实时窗口聚合、实时画像);
  • 一致性保障:通过 Checkpoint + 两阶段提交实现 Exactly-once。

3.2 必须掌握的几个核心概念

  1. 时间语义

    • Processing Time:算子执行该条数据时的机器时间;
    • Event Time:事件真实发生的时间(需要数据里有时间字段);
    • Watermark(水位线):用来标记「某个时间点之前的事件基本都到了」。
  2. 窗口(Window)类型

    • 滚动窗口(Tumbling Window):如每 5 分钟滚动统计;
    • 滑动窗口(Sliding Window):如每 1 分钟滑动统计过去 5 分钟;
    • 会话窗口(Session Window):基于「空闲间隔」划分会话,如用户会话行为。
  3. 状态(State)与容错

    • 键控状态(Keyed State)、算子状态(Operator State);
    • Checkpoint:周期性保存某个时刻所有算子的状态快照;
    • Savepoint:人为触发,常用于升级、迁移、回滚。

3.3 Kafka + Flink 的典型集成方式

  • Source:Flink 通过 Kafka Connector 消费 Topic;
  • 算子链(Operators):过滤、清洗、聚合、连接、富函数、自定义逻辑等;
  • Sink:把结果写入下游:
    • HBase / Redis:实时查询;
    • Elasticsearch / ClickHouse:实时检索与分析;
    • Kafka:继续派发到下游系统。

四、一条典型实时计算链路长什么样?

以「实时用户行为看板」为例,典型链路可以这样设计:

  1. 数据产生

    • 前端 / App 埋点,上报「PV、UV、点击事件、曝光事件」;
    • Nginx / 应用服务记录访问日志。
  2. 数据采集与进入 Kafka

    • 通过埋点 SDK、Fluentd、Filebeat 等采集日志;
    • 写入 Kafka 的 user-behavior Topic。
  3. Flink 实时处理

    • Source:消费 user-behavior;
    • 清洗:过滤脏数据、补充字段、格式统一;
    • 维度关联:如补充用户归属、渠道信息(可以通过维表 join);
    • 窗口聚合:
      • 每分钟实时统计 PV/UV;
      • 每 5 分钟统计各渠道转化率;
      • 对异常波动设置告警规则。
  4. 结果写入与展示

    • 写入 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 规划清单

在你准备上一个实时计算项目时,可以按这份清单过一遍:

  1. 业务目标是否真的需要实时?

    • 需要秒级 / 分钟级;
    • T+1 足够?能否先用离线方案验证价值?
  2. 数据从哪里来?落到哪些 Topic?

    • 是否已经有统一的埋点 / 日志体系?
    • Topic 命名与分区、保留策略是否规划好?
  3. 实时作业要算什么?输出给谁?

    • 指标、维度、窗口类型;
    • 下游系统是看板、风控、推荐还是告警?
  4. Flink Job 如何拆分?

    • 是否拆成多个独立作业,降低耦合;
    • 是否有公共逻辑需要抽成库 / 模板?
  5. 状态、Checkpoint 与容错策略?

    • 状态是否会持续增长?如何清理?
    • Checkpoint 周期、超时时间、存储位置?
  6. 监控、告警与排错路径?

    • 需要哪些关键指标;
    • 事故发生后,谁负责排查、该看哪些图表?

八、总结:Kafka + Flink 是实时计算的「黄金组合」

这一站我们从架构师视角,把 Kafka + Flink 组合在实时链路中的角色和设计要点串了起来:

  • Kafka 提供 高吞吐、可扩展、持久化的实时数据管道;
  • Flink 提供 低延迟、有状态、支持复杂时间语义的实时计算能力;
  • 两者组合起来,能支撑从实时看板、实时风控,到实时推荐、实时监控的一整套业务场景。

关键不在于「我会多少 API」,而在于:
能不能设计出一条稳定、可演进、遇到问题好排查的实时链路。


下一站预告

实时计算链路搭好之后,数据只是「跑起来」了,还谈不上「用得好」。

下一站,我们会把视角从「链路」切换到「治理与协同」:
聊聊 数据治理、元数据管理、数据中台 这些经常被提起却容易被误解的概念,
以及它们与一线开发、架构工作之间到底是什么关系。

敬请期待。

更多推荐