ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Flink流处理架构演进:从批处理到流批一体的技术实践

Flink流处理架构演进:从批处理到流批一体的技术实践 写这篇东西的时候我正盯着Flink UI上那条跑得飞快的实时链路突然意识到一个问题从当年用Storm做实时计算到后来Spark Streaming的微批次方案再到现在Flink几乎成为流处理的事实标准这套架构的演进路径背后是整个大数据领域对实时性理解的蝶变。今天就拿Flink流处理架构的演进脉络来做一个深度梳理把我这几年的实践经验和踩过的坑都翻出来晒一晒希望能给正在研究这个方向的朋友省下一些时间。1. 从批处理到流处理架构演进的时代背景1.1 为什么流处理成了大数据领域的刚需早期的大数据场景核心是解决数据多、算得慢的问题。Hadoop生态的MapReduce也好离线数仓的T1调度也好本质上是把数据攒起来等到夜深人静的时候再集中处理。但现在的情况完全不同了——你打开一个App首页推荐要实时更新你刷短视频点赞量要秒级跳动你下一笔订单风控系统要在毫秒级判断这笔交易是否有风险。数据产生的时刻就是需要被处理的时刻这就是流处理存在的意义。从业务侧看需求流处理解决的不只是快的问题更重要的是改变了数据处理的时间维度。批处理看的是过去完成时流处理看的是现在进行时。比如双11大屏上的实时成交量、网约车的实时订单分配、金融系统的实时反欺诈这些场景如果延迟个几分钟业务价值就已经大打折扣了。1.2 Flink出现之前流处理架构的痛点在Flink流行起来之前主流的流处理方案各有各的尴尬。Storm是最早被广泛使用的实时计算框架延迟确实低但它的短板也很明显只提供了低层次的编程API状态管理这块几乎是半成品你想实现一个精确一次的计算语义得费老大的劲。Spark Streaming用的是微批次方案把实时数据流切成一个个小批次来处理吞吐量上去了但本质上还是批处理思维延迟在秒级而且对事件时间、会话窗口的支持比较吃力。这里就要说到流处理架构演进中的核心痛点到底怎样在保证低延迟的前提下让流处理能够支持有状态的复杂计算。在Flink出现之前业界普遍觉得有状态流处理是一件非常棘手的事——状态存哪里、怎么容错、怎么恢复、怎么保证精确一次语义这些问题始终没有一个优雅的解决方案。Flink架构的出现相当于是重新定义了流处理的标准姿势。2. Flink流处理架构的核心设计拆解2.1 架构分层与核心组件Flink的架构设计有一条非常清晰的逻辑主线就是把所有计算都看成是流。批处理只是流处理的一个特例——一个有界的数据流。这个理念直接决定了Flink的架构分层方式。从下往上分部署层这是Flink架构最底层的基础设施层支持Standalone、YARN、Kubernetes、Mesos等多种部署方式。部署层的设计决定了Flink能否在不同资源管理框架下弹性伸缩。我在实际生产中主要用YARN部署Kubernetes也是大趋势。运行时层这是Flink架构的心脏也就是分布式流处理引擎本身。它负责Task的调度、状态的存储与管理、故障恢复、检查点机制等核心能力。作业提交到集群后这个层帮你处理了所有分布式系统的复杂性。API与开发层包括DataStream API、DataSet API、Table API和SQL。DataStream处理无界流和有界流DataSet主要处理批数据Table和SQL提供了声明式的开发体验。这个分层意味着从底层精确控制到上层快速开发你可以按需选择。扩展库层包括复杂事件处理CEP、Stateful Functions等再往上就是生态系统连接器Connector、Flink CDC等都在这个层面。整个架构的核心思想可以用一句话概括状态驱动的流式执行引擎。Flink引擎本身不区分流和批它完全以流模式运行系统运行时的任务调度、容错机制、资源管理都基于流式处理设计。2.2 关键技术点深度解析真正把Flink架构区别于其他流处理框架的是这几个核心技术点。检查点机制与状态一致性Flink的容错机制核心就是Checkpoint它基于Chandy-Lamport分布式快照算法通过定期生成快照来记录整个作业的运行状态。这里有个概念必须搞清楚Barrier机制。上游算子往数据流里注入Barrier每个算子收到Barrier后会对当前状态做快照所有算子快照完成就完成了一次全局状态快照。通过对齐BarrierFlink可以实现Exact-once语义的容错。我在实际调优中的体会是Checkpoint的间隔时间要有一个合理的平衡。设得太大故障恢复时要回放的窗口就长恢复时间也长设得太小频繁做快照又会消耗大量I/O资源。生产环境一般建议30秒到5分钟之间结合数据量和状态大小来调整。这里有个重要配置项是execution.checkpointing.interval还有一个是execution.checkpointing.min-pause这个参数用来控制两个Checkpoint之间的最小间隔避免Checkpoint持续抢占资源。时间语义与窗口机制Flink支持三种时间Event Time、Ingestion Time、Processing Time。这里有一个陷阱很多刚上手的朋友直接用Processing Time这样处理起来最简单但碰到数据乱序、网络延迟结果就不对了。生产环境对准确度要求高的场景务必要用Event Time配合Watermark机制来处理乱序问题。Watermark的理解有个很形象的类比它就像数据流中的一个标尺表示早于这个时间戳的数据不会再有新数据到来了。比如说你设置Watermark延迟5秒钟就是允许最多5秒的乱序数据。设置得过短会有大量迟到的数据来不及处理设置得过长会延迟窗口计算的时间这个参数需要在业务容忍度和实时性之间找平衡。窗口类型上Tumbling Window、Sliding Window、Session Window各自有使用场景。我遇到最多的坑是Sliding Window的滑动步长过小导致窗口重叠过多、计算压力大增这种问题在设计阶段就要想清楚不要等到上线了再调。状态管理与存储状态是流处理架构里最有价值也最难搞的部分之一。Flink的状态分为Keyed State和Operator State两大类。Keyed State又包括ValueState、ListState、MapState、ReducingState、AggregatingState五种类型。存储后端支持MemoryStateBackend、FsStateBackend、RocksDBStateBackend。生产环境强烈建议使用RocksDBStateBackend因为它在状态量超大时可以落盘不会因为状态过大把整个TaskManager的堆内存打爆。但RocksDB的问题是序列化和反序列化的开销算子的处理性能会有一定下降。一个优化的思路是让状态数据结构尽量扁平化避免过深的嵌套减少序列化开销。任务调度与资源管理Flink把作业拆分成多个Task每个Task并行执行在TaskManager的Slot上。调度策略从早期的Eager Scheduling演进到现在的基于PipelineRegion的Lazy Scheduling后者会分批调度任务好处是能更快发现资源不足的问题避免等所有任务分配完才发现资源不够导致全部失败。这里有个排错经验如果你看到作业提交后长时间处于RUNNING但没有任何Task启动大概率是调度器在等待足够的Slot资源这时候去检查taskmanager.numberOfTaskSlots是什么配置可以拼出来的答案。3. 从1.x到2.xFlink架构演进的关键节点3.1 1.x时代的架构特点与局限Flink 1.x时代整个架构完成了一次从核心引擎到完整平台的关键跨越。1.5版本引入了Native Kubernetes集成1.9版本开始把DataStream和DataSet两套API进行统一尝试。架构的核心引擎基本稳定但有几个明显的短板。第一是批流分离带来的开发割裂。你要写一套实时任务用DataStream API写一套离线任务用DataSet API两套API的语义、算子、调优参数都不通用业务逻辑可能是一样的但代码却要维护两份。第二是Table/SQL在1.x早期版本上的表现还很初级。很多JOIN、聚合操作优化得很粗糙流式SQL的语法支持也比较有限。我印象很深刻当时用Flink SQL做流式维表JOIN得自己实现异步I/O还经常遇到状态膨胀的问题。第三是在资源管理和弹性方面相对薄弱。TaskManager的Slot在作业运行期间是固定的动态扩缩容要靠人工介入或者依赖YARN的资源调度自动化的路径不够平滑。3.2 2.x时代的架构革新Flink 2.0在架构层面做了几件大事彻底改变了Flink的面貌。批流一体化的真正落地2.0版本把DataSet API合并进DataStream API数据集成了统一的流处理接口。无论有界流还是无界流都统一用DataStream来表达。架构层面批处理和流处理共用同一套执行引擎。对开发者来说最大的变化是你可以用一套代码处理批和流开发效率是实打实的提升。对架构而言去掉DataSet API意味着运行时不用维护两套执行语义引擎的复杂度和维护成本都降低了。** Adaptive Execution与动态资源调整**这个演进方向非常关键。Adaptive Execution允许Flink在作业执行过程中根据实际的数据分布、负载情况动态调整并行度。我记得在1.17左右版本已经能看到一些雏形2.x把它变成了生产可用的能力。比如数据倾斜的Scene自适应执行可以根据实际流量自动增缩并行度这个能力在数据量波动剧烈的场景下价值很大。流式数仓架构的提出2.x时代的Flink架构演进不再只是引擎内的事情它开始影响整个数仓架构的形态。Flink Streaming Warehouse可以直接对存储在Iceberg、Hudi等数据湖上的数据做实时的流式读取和写入同时保证事务性和精确一次语义。这种架构下湖仓一体成为可能实时数仓和离线数仓的边界被打破。架构演进的价值评估我这里想多说一句不是新版本的特性越多就越好关键是看你的业务场景能不能吃到这些架构改进的红利。如果你的业务是典型的离线实时双层架构那2.x的批流一体和Streaming Warehouse对你的改造意义就非常巨大。4. 落地实践部署、调优与踩坑记录4.1 集群部署策略与资源规划Flink的部署架构选择直接决定了后续运维的体验。三种主流方式我都有过实际部署经验给大家列个对比。部署方式优势劣势适用场景Standalone部署简单独立集群需自行管理资源HA配置麻烦开发测试环境YARN资源统一管理与Hadoop生态整合好依赖YARN队列调度延迟取决于RM生产环境主流选择Kubernetes弹性伸缩能力强环境一致性高运维门槛高网络存储配置复杂云原生演进目标部署时的资源规划有几个参数需要花心思去算。TaskManager的堆内存和框架内存要分开配置思路简单说就是Frame内存是Flink框架自己用的不要和用户代码的内存混在一起。我通常会按TaskManager总内存的15%左右预留框架内存连接到外部系统的网络缓冲也需要单独考虑。每个TaskManager上的Slot数量是个经典的取舍问题。Slot设置得少每个Task有更多独立资源但并行度上不去吞吐量被限制Slot设置得多并行度高但每个Task的可用资源降低容易产生资源竞争。我的经验是如果您的Task主要是CPU密集型的计算Slot数量设为CPU核数的一半左右比较合适如果是IO密集型的窗口聚合、Join操作可以一个任务占一个Slot。还有一个几乎每个Flink集群都会遇到的问题并行度到底怎么设置。这里提供一个基于吞吐测试的经验式调优方法。先把并行度设成与集群核心总数量相同跑一个基准测试观察CPU利用率和反压情况。如果CPU利用率长时间低于50%说明并行度超出实际负载需要逐步降低如果出现持续反压说明并行度不够逐步增加。4.2 常见问题排查与避坑指南下面这些问题每一个都是我在实战中真金白银踩过坑换来的经验。问题一任务长时间PendingTask一直启动不起来这在YARN部署模式下很常见。第一反应看YARN队列的资源剩余情况。如果队列资源充足看一下Flink作业请求的每个TaskManager内存是不是超过了某个NodeManager能够提供的最大内存值。我遇到过一次Ambari管理的Hadoop集群NodeManager只配置了8G内存我给TaskManager配了10G任务就一直pending。花了一个上午排查网络、依赖、权限最后发现是资源申请就超过了物理限制够折腾。**问题二反压Backpressure反压是流处理系统里最需要优先解决的问题。它本质上是下游处理速度跟不上上游生产速度导致数据积压在TaskManager的输入缓冲中。排查反压有个标准路径先在Web UI的Backpressure标签页看哪些Task处于High状态然后对最接近数据源的High Task进行火焰图分析判断瓶颈是CPU计算密集还是IO等待。如果是频繁的序列化/反序列化导致CPU居高不下考虑优化Row的序列化配置或者调整状态存储方案如果是访问外部存储导致的IO等待优先考虑并行度调整或异步IO。问题三Checkpoint超时或失败Checkpoint失败是最折磨人的问题之一。常见的根源有三个状态过大导致快照写入超时Barrier对齐时间过长导致后续数据阻塞外部存储如HDFS响应抖动。排查时终极大法是把Checkpoint的interval调大、把timeout调大先保住作业稳定运行再逐步缩小找到根因。我遇到最隐蔽的一个问题是RocksDB状态清理策略导致的Checkpoint不一致。社区版本在某些场景下RocksDB的Compaction会影响快照的最终一致性后来的版本通过增加核心状态的预写逻辑解决了这个问题。这也提醒大家生产环境尽量使用稳定的Release版本不要追新不要随意打补丁。问题四数据倾斜数据倾斜是流计算里最经典也最头疼的问题。症状就是某些Subtask的负载特别高其他Subtask空闲整个作业的吞吐被拖住了。排查方法很简单Web UI上看每个Subtask的Records Sent和Bytes Sent如果个别Subtask的数据量远超其他基本就是倾斜了。解决思路有几种对Key进行加盐处理把热点Key打散使用Flink的rebalance或rescale算子重新均匀分布使用CoGroup时注意数据分布的合理性窗口内部可以先做局部聚合再做全局聚合。5. 生态扩展连接器、CDC与血缘管理5.1 JDBC连接器与Sink异常排查Flink的JDBC连接器是使用频率非常高的一个连接器但也是问题聚集地。最常见的异常就是连接断开或连接池爆掉。JDBC Sink写入性能调优的几个关键参数sink.buffer-flush.max-rows控制攒批行数sink.buffer-flush.interval控制刷新间隔sink.max-retries控制最大重试次数。生产环境最常见的错误是连接数打满因为每个并行Subtask默认都会创建自己的数据库连接并行度一大数据库端连接池直接崩溃。排查这类问题第一步看数据库端的最大连接数配置第二步看Flink作业的并行度算一下最多可能建立多少连接第三步给DB连接池配置合理的超时时间和回收策略。还有一个经常被忽略的点JDBC Sink在高并发下频繁写入会导致数据库锁等待可以考虑把batch-size适当调大减少数据库交互次数。5.2 Flink CDC Pipeline部署实践Flink CDC在架构演进中扮演了一个非常重要的角色它让数据同步和实时数据集成变得异常流畅。CDC Pipeline的概念要好好理解它是一整套从数据源到数据目标的数据同步管道支持整库同步、表结构变更自动同步。部署CDC Pipeline时最重要的经验是不要小看数据库Binlog的占用情况。CDC通过读取Binlog来捕获变更如果服务器的Binlog保存时长太短上游稍有延迟Binlog就被清理了CDC任务就断了。我在实际项目里遇到过主库Binlog只有2天的情况刚好某个任务因为升级停了一天多恢复后无论如何追不上新数据只能重新做全量同步再续增量整个过程非常痛苦。建议生产环境结合数据量评估把Binlog保留时间至少设置为3天以上。CDC的任务恢复策略也必须提前想好checkpoint后恢复时Flink CDC会从最近一次checkpoint记录的Binlog位置继续读取。如果变更事件的保留时间足够恢复过程对业务无感如果Binlog刚好在故障窗口内被清理就得做全量重刷。全量重刷期间建议先停止业务写入不然会丢增量。5.3 血缘关系获取与数据治理流处理架构演进到一定规模后血缘关系就成了数据治理的核心需求。数据血缘能解决什么实际问题比如一个下游报表数据异常你怎么快速定位到是哪张源表的数据出了问题哪个Flink任务在产T这些数据血缘关系就是打通这些问题的地图。OpenMetadata是一个开源的数据目录与治理工具可以通过解析Flink作业的元数据来获取作业的输入表和输出表信息进而构建数据集之间的血缘关系图。实际操作中有一个细节Flink作业的元数据信息获取依赖作业中使用的Table名称和Catalog配置是否规范。如果你在代码里硬编码了一些临时表名血缘关系就会很混乱。我在落地血缘关系时有几个经验第一统一管理Catalog与元数据。不要让每张表都是本地临时表尽量让表都纳入元数据中心统一管理血缘关系才能完整覆盖。第二血缘信息要定期刷新。Flink作业的SQL是不断演进的表结构也在变化血缘关系必须跟随作业版本同步更新不能只做一次就不管了。第三血缘关系的价值不仅在追溯还在影响分析。知道某张表要被下线可以通过血缘关系找到所有下游任务和报表提前做变更评估避免半夜突然告警说报表取不到数。6. 流批一体架构落地从选型到实战6.1 如何评估是否需要流批一体架构流批一体是Flink架构演进的大方向但并不是所有团队都需要立刻上。我的建议是先看痛点。如果你们的数仓架构是离线数仓处理T1业务实时计算处理秒级业务两套链路维护成本已经很高而且出现了同一份数据在实时和离线两条链路上计算出来的口径对不上那流批一体就是你们的刚需。如果只是纯实时业务并没有拥抱数仓的场景强行上流批一体反而增加架构复杂度。选型的原则是让架构匹配业务而不是让业务迁就架构。6.2 流批一体架构落地中的关键实践流批一体落地的第一步是数据模型统一。同一份业务数据实时链路和离线链路不能再各自定义一套数据模型。建议从数仓的分层模型出发DWD层的明细数据同时支撑实时和离线两套计算通过Flink SQL同时输出到Kafka和Hive。第二步要处理的是数据回放问题。实时链路处理过程中如果逻辑有Bug需要修复不能只在实时链路打补丁离线链路的数据也需要重放。Flink在流批一体架构下的处理思路是把同一个Flink SQL作业以批的方式提交重新处理历史数据。这里需要注意如果你的作业依赖了事件时间窗口且使用Watermark在批模式下事件时间的处理逻辑和流模式下是不同的需要提前验证边界情况。第三步是保证数据一致性。这依赖于状态存储的选择和Checkpoint机制的可靠性。生产环境的经验是实时链路用RocksDBStateBackend因为状态大是最常见的问题离线重跑时不做Checkpoint追求极致的处理速度。6.3 流式数仓的架构蓝图基于Flink的流式数仓我理解的最佳实践架构大致是这样的数据源通过Flink CDC实时采集进入Kafka作为数据总线Flink做实时清洗和加工实时结果写入StarRocks、ClickHouse等OLAP引擎明细数据落到Hudi或Iceberg数据湖下游通过Presto或Spark SQL查询数据湖兼顾了实时性和批量分析的诉求。这套架构的好处在于实时报表和离线报表的数据底座是同一份口径天然对齐不需要再做一次对账。而且数据湖中的实时明细数据还能支持更长时间跨度的分析和回刷。实际操作中云端资源评估要提前算清楚。比如你要保留7天的实时明细每天的数据量是500GB那就是3.5TB的数据湖存储成本再加上小文件对查询性能的影响。这个成本要在设计阶段就想明白不能等到业务跑起来才纠结。7. 演进中的思考从Flink到整个大数据生态我在Flink架构实践的这几年里有一个越来越强烈的感受技术架构的演进从来都不是孤立发生的。Flink的流处理架构演进实际上是整个大数据生态从存储为中心向计算为中心转变的缩影。以前我们设计架构第一反应是数据往哪里放HDFS还是Hive表然后才考虑怎么算。现在反过来先想清楚这个业务需要什么时效性的计算结果再决定用什么样的计算引擎存储只是计算的延伸。这种思维方式的转变才是流处理架构演进真正带来的深层价值。另外想说的是Flink虽然在流处理领域占据明显优势但也不要神化它。不是所有问题都必须用Flink来解决。批处理的成本效益在特定场景仍然很高特别是大规模历史数据回溯、复杂多表关联分析这些场景Spark仍然是极其优秀的工具。技术选型的正确姿势是根据业务特点、数据体量、时效要求、运维能力综合判断。我个人在实际操作中还有一个很深的体会架构演进最大的阻力往往不是技术本身而是团队认知的迭代。Flink的架构理念需要时间去消化批处理思维和流处理思维的切换在SQL开发、调试方式、问题排查上都有巨大差异。我自己带团队时常用的方法是先把一条核心业务链路从离线改造为实时用实际业务价值说服大家而不是直接大规模铺开。最后再分享一个我在Flink生产运维中始终坚持的小习惯每一次作业上线前都先写清楚这个作业的状态规模预估、Checkpoint间隔设计、重启策略配置、反压告警阈值。流处理作业和批处理作业不一样它一旦跑起来就是7×24小时持续运转前期设计和事前预案的价值远超事后排查。哪怕是一个状态增长过快的小问题如果没在前期管住运行一个月后都可能变成压垮整个集群的大故障。
返回列表