关于流批一体架构的思考

关于流批一体架构的思考 1. 流批一体平台的主要作用数据采集从各种异构数据源如数据库、文件、消息队列等中实时抽取数据包括增量和全量抽取。数据转换对抽取的数据进行清洗、转换和处理以满足目标系统的要求。它提供了一系列的内置转换器和处理器也支持自定义的转换逻辑。数据加载将转换后的数据加载到目标系统中支持多种数据目的地如搜索引擎、推荐引擎、大数据表等。实时监控和报警提供实时监控和报警功能可以监控数据流的健康状态、数据延迟和错误。可视化界面通过可视化界面用户可以轻松配置和管理数据流并实时查看数据流的运行状态和性能指标。2. 数据架构演进2.1 BI 架构BIBusiness Intelligence商业智能的概念很早就有了正如 AI 这一概念一样。早期它的内涵相对模糊按照百度百科的解释“商业智能描述了一系列的概念和方法通过应用基于事实的支持系统来辅助商业决策的制定。” 随着人们实践不断深入BI 系统的样貌也逐渐清晰。到了上世纪九十年代BI 系统迎来了它的第一个辉煌时期Gartner 将各种类型的类 BI 系统全部统称为 BIBI 产品也基本确定为了是一套集数据清洗、数据分析、数据挖掘、报表展示等功能于一体的完整解决方案数据仓库也基于此建立。那时虽然没有大数据的概念但数据分析、商业分析显然是人们长久以来都有的需求也积累了相当多的方法论。当数据量不是主要矛盾时BI 系统能够支持的分析方法、UI 等层面就成为了核心竞争力。当然大多数 BI 系统都构建在关系型数据库之上或者说很多 BI 系统本就是商业关系型数据库的配套产品因此也都是支持 SQL 语言的。以 BI 系统为核心的数据架构如下图所示初代 BI 系统没落的原因主要是底层构建在传统关系型数据库之上因为存在数据一致性约束等问题支持不了大数据。不支持非结构化数据。2.2 传统大数据架构为了解决初代 BI 系统存在的上述问题一些公司开始研发分布式的计算引擎和分布式的存储平台。其中最成功、最知名的便是 Google 研发的分布式文件系统与 MapReduce 计算引擎后来这套技术被开源重写为了 Hadoop 体系的多个项目其生态圈也不断扩大。典型的传统大数据架构流程传统大数据架构它的业务系统数据源可能是关系型数据库 MySQL也可能是平面文件也可能是任意未知的源数据采集和数据同步工具也是视具体的业务和上下游技术选型而定接下来数据会进入数据仓库大致上会依次经过 ODS 层、DWD 层和 ADS 层最终提供给消费方使用。集团内通常业务数据通过 binlog 同步到 TT或者流量日志直接上报到日志服务器再同步到 TT。TT 定期将一个时间区间内的数据同步到 ODPSODPS 再通过每日调度的任务对这些数据进行处理最终落到 ADS 层的表。结果表的数据再同步到 Holo 或 Lindorm 等介质中供消费方使用。因此单看这整个流程实际上就是典型的传统大数据架构的一种实现。但需要注意的是该架构并没有对输入数据有结构化的要求也没有规定 ETL 过程使用的工具和编程语言。在这种架构下业务系统和分析系统的隔离性做得更好了而且无论输入数据是什么最终提供给消费方的都是标准的结构化数据。它的缺点是整个过程不再有完整的解决方案需要做大量的定制化工作。2.3 流式架构流式架构的思路相当激进。虽然传统大数据架构在技术选型上与 BI 系统比已经算是脱胎换骨但其精神还是一脉相承。流式架构干脆扔掉一整套离线的数据采集、数据同步和 ETL 工作直接让流式计算引擎消费业务数据库产生的增量数据并直接输出给消费方以此提供实时的计算结果。而早期的技术储备明显不足以同时高质量保证实时性和结果的准确性因此只被用在了极少数对结果实时性十分敏感却对准确性要求不高的场景中。随着技术的进步和业务复杂度的提高这种架构也基本销声匿迹了。下图是流式架构的典型代表2.4 lambda 架构在早期技术无法同时支持结果的实时性和准确性的情况下如何通过架构设计同时满足实时性与准确性需求Nathan Marz 提出了 Lambda 架构。先看lambda架构的示意图Lambda架构的逻辑是流任务与批任务读取相同的数据源实时计算结果由流任务产出批任务通常按天执行计算T-1的数据并写入到结果表中。最终数据应用根据自己的需要对两个结果表的结果进行合并。其核心思路是:用流任务保证结果的实时性同时用批任务保证结果的最终一致性。Lambda 架构包含离线和实时两条链路两条链路从同一个业务数据源获取数据。离线链路通过定期调度周期性的同步数据源并将其持久化存储在分布式文件系统随后这些数据会被提交给 Spark、MaxCompute 等离线计算引擎进行大规模批处理作业经过离线处理后的结果数据通常会存入高可用、低成本且支持复杂查询的下游存储系统例如 ADBHologres 等通过暴露离线数据的数据服务进而为各类业务场景提供基于历史全量数据的决策支持整个链路提供海量数据的高吞吐、高稳定性的处理能力但结果从采集到产出至用户可见往往需要数个小时。实时链路则通过实时捕捉数据源的 CDC 信息能够近乎实时地追踪并获取数据源的变化情况变更事件一旦产生就会被立即推送到消息中间件随后经过 FlinkSpark Streaming 等流处理引擎消费这些消息能够做到对数据源变更低延迟高容错的秒级响应能力。但 Lambda 架构有几个显而易见的缺点需要开发、维护两套系统成本太大。两套系统难以保证计算口径的一致。不同引擎间由于支持的函数不同、参数配置不同导致代码不能复用容易出现数据一致性和质量问题总之Lambda 架构在满足了部分业务需求的同时给开发和运维同学也带来了 “深重的灾难”。他分拆了计算链路和存储链路导致一个业务逻辑需要维护两套代码不同引擎间由于支持的函数不同、参数配置不同导致代码不能复用容易出现数据一致性和质量问题导致开发成本大运维成本高。个人理解Lambda 就是一种硬把批和流杂糅在一起的架构。2.5 Kappa 架构在流处理技术不成熟的时期主要问题之一就是吞吐量上不去。随着 Kafka 等大数据消息队列的出现吞吐量不再是瓶颈。Kappa 架构的主要贡献之一就是引入了分布式消息队列。如下图所示与 Lambda 架构不同Kappa 架构只保留了流处理层完全舍弃了批处理层。Kappa 架构专注于事件驱动和实时处理只使用一套实时任务链路完成业务数据产出有效降低了系统的复杂性和维护成本同时也提高了数据处理的灵活性和响应速度。但由于实时链路中普遍存在的数据延迟乱序等问题导致计算结果偏差需要业务对当日数据的误差具有一定容忍性。此外数据回刷时需要将大规模数据集重放至消息中间件中这一过程重资源消耗且对消息中间件的存储量和吞吐量有很高要求。对于一些月度或季度汇总的统计周期较长的指标纯实时计算的处理方式会产生大量的排序消耗和中间状态存储。简单来说Kappa 架构让其中一个流处理层正常运行数据应用读取它的输出当数据出现错误或是业务逻辑发生变更时启动另一个流处理层利用消息队列的重播机制重新消费先前的数据并输出到另一个结果表中当确定可以替换线上表时完成替换。总的来说Kappa 架构虽然解决了 Lambda 架构一套逻辑两份代码的问题但同时引入了资源消耗过大、数据存在误差等问题在实践上仍然存在很大的挑战。不过Kappa 架构的另一贡献是启发了人们用单一系统去实现曾经需要两套系统才能实现的需求。人们开始思考为什么流式计算引擎不能提供结果的准确性是哪些环节出了问题流处理引擎是否可以批处理引擎等价的语义2.6 流批一体架构2.6.1 SparkSpark 针对 MapReduce 数据处理过程中海量中间结果落盘影响吞吐量的问题提出了基于内存的计算思想并且通过 DAG 执行引擎优化执行流程宣称能够提高 Hadoop 中 MapReduce 任务效率 100 倍。并且也提供了丰富的 API 支持图计算和机器学习计算功能。Spark 在提出时专注于批任务的处理为了迎合流任务处理的需求Spark 建设了 Spark Streaming 组件能够支持准实时的数据处理。但 Spark Streaming 组件的实现方式是基于微批处理将数据流划分为一系列时间戳连续的数据块每个数据块经过 DAG 执行引擎执行后产生结果。由于 Spark 将实时流切分为微批进行处理实时任务和离线任务仅有批窗口大小的差异实时任务的窗口放大至天级别便是离线任务。因此这种实现方式天然的可以支持一份代码在实时任务和离线任务运行即刚才提到的流批一体。但这种做法只能支持小体量且延时要求不高的应用并且不支持基于业务时间的窗口只能支持数据到达时间的滚动窗口导致很多具有 session 概念的业务数据无法计算。2.6.2 FlinkSpark 将数据流视为特殊的批数据采用微批来处理数据流的做法借鉴了离线数据处理的思想能够保证数据吞吐量但对较低响应延迟的场景却无能为力。Flink 专注于实时任务的处理将批处理视为特殊的有界数据流提供低延迟精确一次基于事件时间的窗口的能力来满足秒级甚至毫秒级的实时数据处理需求。Flink 提出 State 的概念通过维护实时任务的中间状态以及回撤流机制保证结果的准确性来实现单记录粒度的数据处理能力来达到实时响应的要求并通过 checkpoint 提供容错机制最后实现了 Window 和 Time 功能来实现基于事件时间的开窗能力。然而在实际落地过程中用户反馈呈现出明显分化在流处理场景中Flink 凭借强大的状态管理、exactly-once 语义保障以及低延迟性能已成为行业事实上的标准但在批处理场景中许多用户仍倾向于使用 Spark 或 ODPS 等系统因其具备更成熟的查询优化器、更高的吞吐效率以及更友好的开发调试体验。更为关键的是即便基于 Flink 实现流批统一开发者往往仍需为流和批分别配置执行环境、调度策略、Checkpoint 参数及并发度等设置。SQL 层面的统一同样面临挑战Flink 的流式 SQL 引入了事件时间、处理时间、窗口机制等概念而这些在传统批处理场景中往往无需关注。尤其在高吞吐实时场景下系统性能高度依赖对状态清理和反压处理等机制的深入理解要求开发者掌握复杂的调优技巧。这导致 Flink 应用的整体开发门槛居高不下实际落地中仍普遍依赖专业化的实时计算团队支撑。即便业务逻辑可以通过同一段 SQL 表达工程实践中却常常需要维护两套独立的运行模式。更进一步地为了实现 “全量 增量” 的处理流程开发者仍容易陷入经典的 Lambda 架构困境先编写一个 Flink 批作业处理历史数据以生成基线再另起一个流作业消费实时日志进行增量更新。尽管两者的处理逻辑高度相似但由于执行模式有界 vs 无界、资源配置和运维方式的不同不得不重复开发、分别部署和独立维护极易因版本错配或参数差异导致逻辑偏差进而引发数据不一致问题。这种困境的背后是流与批在语义层面的深层差异批处理通常基于静态快照进行一次性计算而流处理则依托事件时间、窗口触发与持续更新机制具有动态性和状态演化特性。在迟到数据处理、重复记录去重、聚合结果修正等场景中两者的处理逻辑天然不同。即使使用同一段 SQL若未严格约定时间语义、窗口策略与聚合行为仍可能产出不一致的结果。因此尽管 Flink 在运行时层面实现了流与批执行模型的统一显著推进了架构收敛但距离开发者所期待的 “一套逻辑、一次开发、全域生效” 的理想目标仍有不小的差距。2.6.3 Flink、Spark和阿里巴巴流批一体平台SARO对比维度FlinkSparkSaro典型场景实时风控、IoT 边缘计算、实时大屏、金融级一致性要求场景离线数仓 ETL、机器学习训练、报表统计等高吞吐批场景中大型企业需快速构建流批一体链路的业务团队如电商价格 / 库存域核心思想将批视为有界流的特例把流 “切片” 成小批量再复用批处理引擎执行本质是 “用批模拟流”并非底层计算引擎而是基于 Flink 构建的流批一体数据开发平台。它不替代 Flink而是对其能力进行封装与增强通过 DSL 抽象、可视化编排、自动任务生成批 流双作业等方式降低流批一体落地门槛。优势低延迟毫秒级、强状态管理、事件时间语义完备天然支持事件时间Event Time、乱序处理、精确一次Exactly-once语义及低延迟状态管理生态成熟、SQL 优化强大、社区活跃、批处理性能极致开发效率高拖拽 DSL、运维标准化、支持热点隔离与成本优化