架构

以业界最主流的 Kafka + Flink 组合为例,架构分为四层,是实时大屏、实时风控、实时推荐等业务的标准落地方案。

结果输出与应用层

流计算引擎层 Flink

Kafka 消息缓冲层

数据接入层

业务系统
订单/支付/用户事件

日志埋点
应用日志/用户行为

数据库CDC
MySQL Binlog

IoT设备
传感器/终端上报

SDK直接生产

Filebeat/Fluentd采集

Debezium/Canal同步

IoT网关/MQTT Broker

Kafka 集群
多副本 + 分区

业务事件Topic

日志行为Topic

数据变更Topic

设备指标Topic

Flink 分布式计算集群

数据清洗/过滤

多流关联/维表Join

窗口聚合/统计

规则匹配/风控检测

特征工程计算

分析存储
ClickHouse/Doris/ES

实时大屏/多维分析

业务存储
MySQL/Redis

实时推荐/业务接口

告警通道
钉钉/短信/风控拦截

实时告警/风险拦截

回流Kafka
结果Topic

下游业务二次消费

1. 数据接入层

负责将多源异构数据统一写入Kafka,常见接入方式:

  • 业务事件:后端服务通过SDK直接生产消息(如订单支付、用户注册事件)
  • 日志埋点:通过Filebeat、Fluentd采集应用日志、用户行为埋点
  • 数据库变更:通过Debezium、Canal采集MySQL Binlog,实现CDC数据同步
  • 设备数据:IoT网关、MQTT Broker批量上报终端指标数据

2. Kafka 消息缓冲层

承上启下,解耦数据源与计算引擎,同时提供数据缓存与重放能力:

  • 按业务域拆分Topic,如 trade_order_topic 、 user_behavior_topic
  • 分区数根据峰值吞吐设计,通常与下游Flink算子并行度对齐
  • 配置多副本保障高可用,数据保留时长覆盖业务容错窗口(通常7天以上)

3. 流计算引擎层

实时计算的核心,主流技术选型:

  • Apache Flink:首选方案,毫秒级延迟、强状态管理、原生支持Exactly-Once,适配绝大多数实时计算场景
  • Kafka Streams:轻量级方案,无需额外部署计算集群,基于Kafka原生API实现简单流处理
  • Spark Structured Streaming:适合批流一体、分钟级延迟的场景

核心能力:数据清洗过滤、多流关联、窗口聚合、规则匹配、特征工程等。

4. 结果输出与应用层

计算结果对接下游存储与业务系统:

  • 分析存储:ClickHouse、Doris、Elasticsearch,支撑实时大屏、多维查询
  • 业务存储:MySQL、Redis,对接业务接口、实时推荐
  • 告警应用:钉钉/短信告警、风控拦截,触发实时业务动作
  • 回流Kafka:计算结果写回新Topic,供下游业务二次消费

落地示例:电商实时GMV大屏

接入层 → 订单支付事件写入Kafka → Flink按类目/地区做滚动窗口聚合 → 结果写入ClickHouse → 可视化大屏秒级刷新

核心最佳实践

一、Kafka 侧设计最佳实践

  1. 分区与并行度对齐
    Flink消费并行度建议与Kafka分区数相等,或成整数倍;一个分区对应一个消费线程,避免资源浪费或消费瓶颈。
  2. 分区键合理设计
    必须按业务主键(用户ID、订单ID等)设置分区Key,保证同主键数据进入同一分区,避免流计算聚合时数据拆分、结果不准。
  3. 结构化序列化协议
    优先使用Protobuf/Avro,配合Schema Registry做字段版本管理;相比JSON体积更小、解析更快,同时避免上下游字段不一致导致的计算异常。
  4. 数据保留留足冗余
    保留时长需覆盖最大故障恢复窗口,比如支持重算3天数据,则至少保留7天,为故障排查、数据补算预留空间。

二、流计算侧(Flink)最佳实践

  1. 端到端 Exactly-Once 保障
  • 开启Flink Checkpoint(建议间隔1~5分钟),统一管理消费位点,不依赖Kafka原生offset提交
  • Sink端支持两阶段提交(如Kafka、JDBC Sink),实现生产端到消费端全链路数据不丢不重
  1. 乱序与迟到数据治理
    通过水位线(Watermark)定义乱序容忍度,配合窗口允许迟到机制;极端迟到数据走侧流输出补发,避免窗口计算结果失真。
  2. 状态与性能优化
  • 大状态场景(天级窗口、海量Key聚合)选用RocksDB状态后端+增量Checkpoint,降低内存压力
  • 热点Key采用两阶段聚合:先本地预聚合打散热点,再做全局聚合,避免单节点负载过高
  1. 消费组隔离
    不同实时业务使用独立消费组,避免任务之间互相影响;同一份数据可被多个消费组重复消费,互不干扰。

三、运维与监控最佳实践

  • 核心监控指标:Kafka侧监控消息堆积量、ISR同步状态、吞吐速率;Flink侧监控Checkpoint成功率、算子反压、消费延迟、状态大小
  • 故障自愈:任务配置自动重启策略,基于Checkpoint恢复位点,无需人工补数
  • 压测前置:上线前按峰值流量的1.5倍做压测,验证Kafka吞吐与Flink计算能力是否达标

更多推荐