尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

从分钟级到毫秒级:基于Flink的实时流处理架构设计与实践

从分钟级到毫秒级:基于Flink的实时流处理架构设计与实践 运营在周会上丢过一句话数据大屏能不能别总慢五分钟我当时手里那套还是T1跑批改造到分钟级已经是极限。后来我把目标直接定在毫秒级整套架构用 Flink 重写端到端延迟从分钟级降到几百毫秒。这篇文章就把这次实时流处理架构设计从需求拆解、组件选型、核心实现到踩坑记录完整拆开讲适合正在做实时数仓、实时风控或想做实时架构升级的开发和架构师。先交代一下背景业务量不算极端但每天也有几十亿条事件要处理峰值分钟级千万级。目标不是做实验室里的 Benchmark而是让业务真实跑起来后从数据产生到业务可查延迟能稳定保持在秒级以内多数场景做到几百毫秒。1. 先理清需求毫秒级到底意味着什么先别急着打开 Flink 官方文档。做架构设计最忌讳拿到需求就直接画拓扑尤其“毫秒级”这个词不同业务场景的语义完全不一样。1.1 从“准实时”到“真实时”的延迟界限在实时计算这个领域大家经常把“准实时”和“真实时”混在一起。我的标准很简单分钟级、秒级延迟比如 Spark Streaming 的微批模式本质是攒一批算一批属于准实时毫秒级延迟事件进来之后很快触发计算属于真实时。Flink 能实现毫秒级核心原因是它的事件驱动架构和流式执行引擎。数据记录一到算子就被处理而不是等到批次边界。这里要区分两个概念单算子处理延迟和端到端延迟。单算子处理延迟常常是毫秒甚至微秒级但端到端延迟还包括消息队列排队、网络传输、状态读写、下游写入耗时所以要谈“毫秒级架构”一定是指整套链路。实际项目里我会把“毫秒级”拆成三个可量化的指标数据产生到进入 Flink source 的时间、Flink 内部处理时间、Flink sink 到下游存储被查询到的时间。三个指标里最容易被忽略的是第一个和第三个它们往往决定端到端最终能不能到秒级以内。实操心得不要只盯 Flink Dashboard 上的某个算子延迟。真正的端到端延迟要在业务数据里埋时间戳从业务发生时刻算起到下游能查到这条数据为止整个链路做一次计时。有一个注意点毫秒级不是所有场景都必要也不是所有场景都能做到。比如下游是每分钟刷一次的报表你再怎么优化 Flink 也白搭。架构设计的第一步是把延迟目标定清楚并且让业务方也认可这个目标。1.2 架构选型前的三个核心问题再做任何一个实时项目我习惯先用一页纸把需求发散开。这里说的“发散”就是标题里那个意思不要一上来就钻进 Flink 的某个算子细节先把整个问题的边界画出来。第一个问题数据源是什么类型。是埋点日志、业务库 binlog、还是消息队列里的已有事件流不同类型决定接入层的选型。日志型数据通常打到 Kafkabinlog 型数据适合用 Flink CDC 直接捕获。第二个问题峰值吞吐和延迟要求。峰值每秒多少条、突发流量倍数是多少、能不能接受偶发背压这些数字决定 Kafka 分区数、Flink 并行度、状态后端选型。我见过不少项目需求只说“数据量不大”结果上线当天就被流量打爆。第三个问题下游怎么消费计算结果。是实时大屏轮询还是在线服务接口查询或者写入数仓做进一步分析下游决定结果存储和查询层设计。比如大屏一般用 WebSocket 推送在线风控一般用 Redis 或内存表服务数据分析场景可能落到 ClickHouse 或 Doris。这三个问题发散完之后再收敛成一张架构分层图。有人喜欢用 TOGAF 那套企业架构框架来画有人习惯自己画一张四层图我的经验是框架不重要重要的是每一层之间的数据流和延迟预算写清楚。例如接入层预算 100ms计算层预算 200ms存储查询层预算 300ms加起来到秒级以内这个预算后面要严格盯。2. 整体架构设计与组件选型把需求理清之后才开始聊组件。我这次架构最终选的是 Kafka Flink Redis/ClickHouse 这套组合下面说清楚每一层为什么这么选。2.1 数据接入层为什么一定要加一层消息队列我先说结论只要业务允许毫秒级架构的接入层一定要放消息队列最常见的是 Kafka。原因不是“大家都在用”而是消息队列解决了两个关键问题。第一个是削峰填谷。业务流量天然有高峰低谷如果让 Flink 直接对接数据源峰值来了要么拒绝数据要么把资源拉满有了 Kafka 缓冲Flink 可以按自己最舒服的速率消费哪怕峰值持续几分钟也不丢数据。第二个是多消费者解耦。同一份数据实时大屏要一份实时风控要一份实时数仓也要一份。如果没有消息队列每个下游都要直接对接数据源数据源压力会成倍增加通过 Kafka 的消费组机制每个下游独立消费互不影响。Kafka 的分区数要提前想好。我一般按目标峰值吞吐除以单分区消费能力来估算单分区写入能力通常在几 MB/s 到几十 MB/s读取要更快一些。分区数同时决定 Flink source 并行度的上限Kafka 有 64 个分区Flink source 并行度最大就设 64。太小的分区数会限制后续扩展太大又浪费资源所以要压着未来半年的业务增量来定。2.2 计算引擎层Flink 核心优势与并行度规划计算引擎我几乎没有犹豫就选了 Flink。Spark Streaming 我之前也用但微批模型在延迟上天然吃亏哪怕把批处理间隔压到 500ms也会因为调度、批次等待带来额外开销。Flink 是真正的流式执行引擎数据一条条处理事件时间支持也更成熟。这里展开说三个设计要点。第一状态管理。实时计算大部分场景不是单纯过滤而是要统计、关联、去重比如统计用户最近 5 分钟加购次数。这需要保存中间状态。Flink 的状态后端我推荐 RocksDB虽然性能比纯内存的堆内存状态后端慢一点但状态可以存储在本地磁盘不会因为状态太大撑爆 JVM 堆内存。状态量大的场景优先 RocksDB状态量小且追求极致性能可以用堆内存。第二Checkpoint 机制。要保证毫秒级延迟又不能丢数据必须开启 Checkpoint。我通常把 checkpoint 间隔设在 30 秒到 1 分钟状态很大时再拉长。间隔太短频繁做快照会影响性能太长故障恢复时丢失的数据会变多。这里没有绝对标准要靠压测调整。第三并行度规划。并行度不是越大越好。并行度太大网络 shuffle 增加状态分布变散反而更慢。我的做法是压测时从 Kafka 分区数出发先让 source 并行度等于 Kafka 分区数然后逐步调大下游算子并行度观察吞吐和背压变化。对比项FlinkSpark Streaming处理模型流式逐条处理微批处理典型延迟毫秒级秒级状态管理原生状态后端支持 RocksDB依赖外部存储较多事件时间支持原生成熟支持但配置复杂注意Flink 的每个算子并行度可以不同。source 并行度受上游分区数限制transform 和 sink 并行度可以独立设置。别图省事一个全局并行度用到底。2.3 存储与查询层结果存储与服务化Flink 算完结果之后存在哪里直接决定了业务方能不能达到毫秒级体验。这部分我最常说的一句话是Flink 负责算得快存储查询层负责查得快两条腿缺一不可。不同类型结果选不同存储。简单计数、去重后的汇总结果适合放 Redis毫秒级查询没有问题明细结果、需要多维聚合分析的放 ClickHouse 或 Doris如果结果要和业务库的维度数据做关联查询TiDB 这类分布式关系库也值得考虑Flink SQL 可以直接通过 JDBC 连接 TiDB 写入或读取TiDB 的实时分析能力近几年进步很明显。需要注意一个常见误区有人会把计算结果直接写到 MySQL 一张大表让业务方去查。数据量小没问题数据量一上来MySQL 的查询延迟就压不住。我自己做过一次对比同样的结果集从 MySQL 查耗时一两秒从 ClickHouse 查只需要几十毫秒这个差距在实时场景下非常致命。大屏场景还有一点别让前端周期轮询后端接口几百毫秒的接口延迟加网络开销用户体感会差很多。更好的方案是 Flink sink 到 Redis 后后端通过发布订阅或者 WebSocket 把增量结果推给前端这样端到端路径更短实时感更强。3. 核心链路实现从 Kafka 到 Flink 再到下游架构层聊清楚了下面进入实操环节。我以一个通用的实时指标统计场景为例从 Kafka 读取用户行为事件统计每个商品最近 5 分钟的加购次数结果同时写入 Redis 和大屏。3.1 环境准备Flink 部署与版本选择先交代版本。我这次用的是 Flink 2.2.1Flink CDC 插件是 3.5.0。这里提醒一句不要盲目追新也不要一直守着老版本。Flink 2.x 相比 1.x 在 SQL 能力、资源管理等模块上变化不小如果团队是从 1.x 迁移过来升级前一定要把用户代码的兼容性测一遍。部署方式上小团队验证可以先用 Standalone 单机模式或者用 Docker 在本地把 Flink 跑起来。只是测试别一上来就上 Kubernetes先让业务逻辑跑通再考虑背后的资源调度和容灾。Docker 部署 Flink 时我踩过的第一个坑是内存限制容器默认堆内存和容器内存上限对不上作业跑着跑着就被 Kill。解决办法是在 docker-compose 或启动命令里显式指定 JVM 堆内存和容器内存比如docker run -it -d \ --name flink-jobmanager \ -m 4g \ -e JOB_MANAGER_MEMORY2048m \ -p 8081:8081 \ flink:2.2.1 jobmanager这里的核心逻辑是-m 4g 限制容器总内存JOB_MANAGER_MEMORY 指定 JVM 进程的堆内存两者之间要留出 JVM 元空间、线程栈和网络缓冲的空间不然容器会频繁触发 OOM Kill。生产环境我更推荐 Flink on Kubernetes 或 Flink on YARN。通过 Kubernetes Operator 管理作业可以做到作业级资源隔离和自动重启。这块内容可以单独写一篇这里不展开。3.2 实现一个毫秒级实时 ETL 作业我习惯用 Flink SQL 来实现大部分实时任务简单直观维护成本低。下面这个例子就是加购指标统计的核心 SQL。CREATE TABLE kafka_behavior ( user_id BIGINT, item_id BIGINT, behavior STRING, event_time TIMESTAMP(3), WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic user-behavior, properties.bootstrap.servers kafka-1:9092,kafka-2:9092, properties.group.id realtime-metric-group, scan.startup.mode latest-offset, format json ); CREATE TABLE redis_sink ( item_id BIGINT, cnt BIGINT, window_end TIMESTAMP(3), PRIMARY KEY (item_id) NOT ENFORCED ) WITH ( connector redis, redis.mode single, redis.host redis-1, redis.port 6379, format json ); INSERT INTO redis_sink SELECT item_id, COUNT(*) AS cnt, TUMBLE_END(event_time, INTERVAL 5 MINUTE) FROM kafka_behavior WHERE behavior cart GROUP BY TUMBLE(event_time, INTERVAL 5 MINUTE), item_id;这段 SQL 里最需要注意的是 WATERMARK 那一行。业务数据是源源不断进来的但网络抖动、客户端重试会导致事件时间乱序比如一条 12:00:01 产生的数据可能 12:00:03 才到达。Watermark 的作用是告诉 Flink这条线之前的事件时间都已经齐了可以触发窗口计算了。我把乱序容忍度设为 5 秒意思是允许事件时间偏差 5 秒内的数据迟到超过 5 秒的不再计入窗口。窗口类型这里用的是 TUMBLE固定 5 分钟滚动窗口。如果你要的是实时大屏上的近 5 分钟指标滚动窗口是最简单的方案因为它每个窗口结束后一次性输出结果。要是想更平滑可以换成 HOP 窗口或者用连续 TopN但复杂度会上升建议从滚动窗口起步。3.3 Watermark 与事件时间乱序数据不影响毫秒级延迟上面 SQL 中已经出现了 Watermark这里单独拎出来讲因为这是实时流处理最容易踩坑的地方热搜里也经常有人搜“Flink SQL 中 water”。很多人一开始会用处理时间PROCTIME代替事件时间因为处理时间不需要关心 Event Time 和 Watermark写起来也简单。但处理时间有个致命问题它表示的是 Flink 机器处理这条数据的时间不是数据真实发生的时间。一旦数据在链路里积压了 30 秒处理时间就比事件时间晚了 30 秒统计出来的结果就不准了。所以在实时统计场景我强烈建议用事件时间。Flink SQL 里定义 Watermark 的语法就一行WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND这行的意思是允许事件时间最多迟到 5 秒。这个 5 秒不是拍脑袋定的要根据上游 Kafka 积压情况和业务容忍度来设置。设得太大窗口要等更久才触发延迟上去设得太小乱序数据被丢掉结果不准。我一般先用线上数据回放统计 P95 的事件时间延迟然后在这个基础上加一点冗余。还有一个细节Flink SQL 里每个 Source 表都可以单独定义 Watermark。如果一个作业 join 多个数据源每个源的事件时间口径要一致不然 join 窗口会合不上数据就莫名其妙少了。4. 毫秒级链路里的典型坑与排查光看文档永远不知道生产环境会出什么幺蛾子。下面这几个问题是我在这套架构里真实踩过的每一个都花了至少半天才定位清楚写出来给大家省点时间。4.1 Flink JDBC 连接器异常Flink 的 JDBC 连接器比如写入 MySQL、TiDB、PostgreSQL是使用频率最高的连接器之一也是最容易出问题的一个环节。最常见的报错有几种连接超时、Too many connections、事务提交失败、writer 线程阻塞。先说 Too many connections。Flink JDBC sink 的每个并行子任务默认会维护自己的连接池比如并行度是 16每个子任务连接池配了 5 个连接那就是 80 个连接。如果下游 MySQL 的 max_connections 设的是 200多个作业一跑连接数直接打满其他服务也跟着遭殃。解决办法有两个层面。作业层面控制并行度和连接池大小比如WITH ( connector jdbc, url jdbc:mysql://tidb-1:4000/app_db, username root, password ******, sink.buffer-flush.max-rows 1000, sink.buffer-flush.interval 1s, sink.max-retries 3 )sink.buffer-flush.interval设为 1 秒让数据攒在一起批量写入能大幅减少连接占用和写入次数。代价是引入最多 1 秒的额外延迟对于毫秒级目标来说一般可以接受。你要纯毫秒级写入也不是不行但代价是连接更多、数据库压力更大这个平衡要自己把控。第二层面是数据库侧把 JDBC 驱动连接池的最大连接数、空闲连接回收时间调好别让连接长期空闲后又超时。TiDB 做实时数仓下游时同样建议用 JDBC sink 批量写入单条 upsert 的性能和批量写入差一个数量级。异常现象常见原因处理建议Too many connections并行度乘连接池过大调小 sink 并行度开启批量写入连接超时网络不通或连接空闲回收时间太长检查网络调小空闲超时重复写入至少一次语义导致下游去重或换精确一次方案注意JDBC sink 默认是至少一次语义断线重连或重试时可能重复写入。如果业务要求精确一次建议把结果先写入消息队列或对象存储最后再事务性写入目标库或者使用支持两阶段提交的 sink 插件。4.2 Flink SQL 与 SQL Gateway 实践Flink SQL 的入口有好几个SQL Client、SQL Gateway、以及 DataStream 中嵌 SQL。很多教程让新手直接改 SQL Client 配置文件生产上其实不推荐这样用。SQL Client 是本地开发调试用的适合你一个人交互式敲 SQL不适合作为团队共享或自动化发布的手段。生产环境我推荐用 SQL Gateway可以理解为一个常驻的 Flink SQL 服务端。把 SQL 提交到 GatewayGateway 负责解析、生成作业、提交到集群这样客户端就不用关心 Flink 集群通信细节。配合作业提交工具可以让不同团队提交 SQL 作业时不直接接触底层集群。用 SQL Gateway 还有一个好处它天然支持 REST API可以很方便地集成到内部的发布平台或 CI/CD 流水线里。比如一个实时数仓团队分析师维护 SQL 文件发布平台调用 Gateway API 提交作业整个过程可以做到版本管理和权限控制。踩过的一个坑是SQL Gateway 默认会缓存 Session 和 Catalog如果你改了 Kafka Topic 映射关系但 Session 没刷新作业看起来提交成功实际读的还是旧表结构。排查方式很简单提交作业后立刻去 Flink Web UI 上看作业的血缘和执行计划确认读写的表和预期一致。4.3 CDC 链路Flink CDC 3.5.0 与 Docker 部署注意点实时链路里很大一部分数据来源于业务库也就是 binlog。我之前用 Canal 把 MySQL binlog 同步到 Kafka再用 Flink 消费链路长、组件多、排查麻烦。后来切到 Flink CDC直接用 Flink 读 binlog省掉了中间那一层。Flink CDC 的部署方式很简单一个 SQL 作业就能建表。比如CREATE TABLE mysql_users ( id INT, name STRING, update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector mysql-cdc, hostname mysql-1, port 3306, username cdc_user, password ******, database-name app_db, table-name users );这里要提醒几个 Docker 部署时容易踩的坑。第一Flink CDC 作业里会包含一个快照读取过程首次启动要全量读一遍数据这时会对源库产生压力建议指定scan.incremental.snapshot.chunk.size控制每个分片的大小别一次性把整张表读出来。第二CDC 连接器需要访问源库的 binlogDocker 网络的宿主映射要配好别让容器内访问不到 MySQL 的 3306 端口。第三Flink 2.2.1 配合 Flink CDC 3.5.0 时要注意 CDC 插件版本和 Flink 版本的兼容性不能用错版本组合否则启动作业时会出现类冲突或找不到 Driver 的报错。我用的那套 Docker 环境大致是 Flink JobManager 和 TaskManager 各一个容器MySQL 一个容器Redis 一个容器再起一个 Flink CDC 作业提交容器。通过 Docker Compose 统一管理方便本地复现和测试。这里不再贴完整 compose 文件网上模板很多照着改端口和镜像版本即可。4.4 Flink SQL 中的 Watermark 常见误用虽然前面写过 Watermark 的基础用法但实际项目里还是看到过不少误用。一种是 Kafka 消息本身带了业务时间戳但建表时没解析出来用 PROCTIME 替代导致凌晨低峰期数据延迟统计崩了另一种是 Watermark 定义在中间临时表上结果 Source 表并没有暴露正确的事件时间字段。检查 Watermark 是否生效最快的方法是在 Flink Web UI 的作业图里看每个算子的 Watermark 值如果始终停在 N 分钟前说明上游数据的时间戳解析有问题或者 Source 表的 WATERMARK 定义字段写错了。可以先用LOCALTIMESTAMP或简单的SELECT *打印几条数据出来确认 event_time 字段确实有值再去看窗口计算。5. 进阶从功能实现到治理与优化当链路能稳定跑起来、延迟也在目标范围之内考验才刚刚开始。实时架构和离线架构最大的不同在于它是 7x24 小时不停转的没人能总盯着作业所以治理和自动化异常重要。5.1 数据血缘与任务治理热搜词里有“Flink 数据血缘”这是很多人忽略的一块。实时作业一变多你可能会遇到这种问题一个 Kafka Topic 被 5 个作业消费改动 Topic 字段格式到底会影响哪些下游如果没有数据血缘记录只能一个个作业去查效率极低。Flink 生态里数据血缘一般靠 Catalog 和作业提交平台来收集。Flink SQL 建表时使用统一的 Catalog 注册元数据提交平台解析 SQL 中的血缘关系写入一张血缘关系表。这样当某个上游字段要变更时可以快速反查受影响的下游作业和字段。我见过做得比较完善的公司会把数据血缘和数据质量联动血缘里发现某个表的产出作业失败就立刻触发下游数据质量规则校验并通知对应用户。这件事不需要一步到位但至少可以从记录开始。5.2 性能压测与资源调优最后聊性能。毫秒级架构不是部署完就有的一定要做压测把延迟真实暴露出来。我自己的压测做法是写一个模拟数据生产者脚本按线上峰值的 1 倍、2 倍、5 倍往 Kafka 打数据同时用 Flink Web UI 的背压监控和业务数据自身的时间戳来评估端到端延迟。如果 2 倍流量时就出现明显背压要么增加并行度要么检查是否有算子处理逻辑过重比如在 Flink 里做了大范围的维表 join这种操作很容易成为瓶颈。几个有效的调优点调大execution.buffer-timeout会让系统攒更多数据再发往下游提升吞吐但增加延迟追求毫秒级延迟时这个值可以调小到 10ms 左右。开启对象重用pipeline.object-reuse减少序列化开销却要小心代码里别持有对象引用。合理设置 TTL状态只保留业务需要的那几十分钟过期数据自动清理减小状态体积。实操心得调优一定是一次只改一个参数改完压一次测记录延迟和吞吐。不要同时动一堆参数出了问题根本没法定位是哪个参数引起的。5.3 从 Flink 发散开这套架构还能用来做什么文章标题里我提到“发散创新”不只是发散需求更是发散应用场景。同一套 Kafka Flink 存储查询的架构换一下数据处理逻辑就能覆盖很多实时业务实时风控Flink 消费用户行为流配合规则引擎做毫秒级规则匹配命中后调用下游接口拦截实时推荐Flink 实时计算用户最近浏览的商品类目更新用户特征供推荐服务查询实时数仓Flink 做 ODS 到 DWD 的实时清洗聚合结果落到 Doris 或 ClickHouse支撑分析师快速查询监控告警Flink 消费机器指标流滑动窗口内指标超过阈值就触发告警。这些场景的底层架构几乎一样区别只在于业务逻辑。所以值得在架构设计阶段多投入一点精力把接入层、计算层、存储层解耦做好后面接新业务就只是写 SQL 写逻辑的事不用重搭框架。我自己在实际项目中感受最深的一点是毫秒级不是靠某一个组件魔法般地实现的而是每一层都做到该做的事不留明显短板。Kafka 做好缓冲Flink 做好计算Redis 或 ClickHouse 做好查询再配合合理的 Watermark 和资源调优最终才能稳定达到目标。如果某个环节延迟特别高其他环节再快也会被拖累。这套思路如果也能帮你在实时流处理架构设计上少走几次弯路那就值了。你手头如果也在做类似的实时链路欢迎按上面的思路先去盘点一下自己的接入、计算、存储三层看看瓶颈到底在哪一层。
返回列表