ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

流批统一实战:从Lambda迁移到Flink SQL与Hudi

流批统一实战:从Lambda迁移到Flink SQL与Hudi 1. 流批统一到底在解决什么问题先说个容易跑偏的事。很多人搜“lambda函数”“lambda表达式 java”跑来看到Lambda架构以为是Java 8语法那套东西。确实同名但完全是两码事。大数据领域的Lambda是一套实时和离线双跑的架构思想名字取自函数式编程的λ演算实际含义是用两条独立的链路分别处理流和批。这里先把这个岔路口说清楚不然后面全是浆糊。Lambda架构和Kappa架构的争论本质上是流批统一这个老大难问题的一个缩影。早期数据平台通常是“T1离线为主实时为辅”白天跑批晚上出报表。后来业务对时效性的要求上来——风控要秒级识别、大促要分钟级指标、推荐要做近实时特征——于是大家开始补实时链路。补着补着就发现离线一套代码、实时一套代码同一个指标在两套作业里算出来经常对不上。你说数据也没错但是口径不一致业务就找上门问到底哪个是对的流批统一要解决的就是这摊事。它不是让你非要选Lambda还是选Kappa也不是让你把所有东西都换成Flink或者Spark。它真正解决的是同一份数据、同一套计算逻辑、同一个口径在流和批两个场景下都能跑运维面要收敛开发成本要可控历史重算要有手段。我做了几年数据平台最大的感受是很多团队不是不知道Lambda和Kappa的区别而是不知道自己业务的真实诉求于是选了最吃力不讨好的组合。比如明明业务只需要十分钟级延迟结果硬上秒级流计算一套链路三个人维护得不偿失又比如业务确实要秒级实时结果偏偏选了Kappa的思路靠流重放补批计算数据量大一点就撑不住。这篇文章不聊抽象概念就聊怎么判断、怎么选、怎么落地。包括我最常被问到的几类问题一次性说清楚。2. 先看明白两种架构的核心逻辑2.1 Lambda的“双跑”设计稳但代价高Lambda架构经典的三层——批处理层、速度层、服务层。批处理层用类似Hive、Spark批处理的方式对全量数据计算产出精确的历史结果速度层用Flink、Storm这类流引擎处理实时增量产出低延迟的近似结果服务层把两者的结果合并对外提供服务批的结果用于校正流的结果用于时效。为什么要这么设计核心原因是早期流计算做不到精确且低成本的重算能力。流引擎处理完就丢状态也不能跨长时间窗口保存。而批处理能做全量重算、能做复杂Join、能保证最终一致性唯一的缺点是慢。既然如此就让快的去快让准的去准各司其职再在查询层合并。这个设计的优点很直接每一层都用自己最擅长的引擎链路稳定边界清晰。当时很多大厂的核心数据平台都是这么搭的。缺点也很明显同一套业务逻辑要写两遍一遍批处理一遍流处理。任何一处口径改动两套代码都要同步修改改漏了实时和离线就对不上。加上数据回流链路长实时结果和批结果中间还隔着周期性的合并任务排查问题的时候经常要查三层日志人累得够呛。所以现在很多团队不是不想淘汰Lambda而是被维护成本拖着走不下来。代码、脚本、调度、监控、依赖关系全部都是双份光梳理就要费不少劲。2.2 Kappa的“重放换重算”简洁但有前提Kappa架构的核心思想是既然流处理引擎能持续跑数据那就把批处理也当成流处理来做。所有历史数据先灌进消息队列比如Kafka用同一个流作业从头消费跑出最终结果。重算时不需要停作业直接另起一个作业从某个历史位点重新消费跑完再切换读取新结果。这听起来很美——一套代码一次开发没有双链路。但Kappa有一个非常硬的前提消息队列必须保留足够长时间的历史数据并且你有能力在合理时间内把所有历史数据再消费一遍。Kafka虽然能通过调整保留策略存比较久的数据但存储成本、broker压力、消费吞吐都会受限制。数据量大到几十T、上百T的时候历史重放不是几分钟的事可能要跑好几个小时甚至按天计算。而且如果流作业过程中依赖了大状态、复杂窗口、外部维表Join重放时状态重建的成本不可忽略。讲个我实际遇到的例子。有个业务想做Kappa化改造把一年的历史订单数据进Kafka准备用流作业全量重算客户标签。Kafka topic的partition调整、压缩配置、磁盘扩容先折腾了一周真正开始消费的时候发现上游数据有几处字段含义历史上变更过同一份JSON在三个月前后结构都不一样。流作业跑到一半直接反序列化失败作业重启又得从最早位置来一遍。最后这个项目砍掉了老老实实继续用批处理补历史、流处理接增量也就是走了Lambda的路。Kappa适用的场景我的判断是数据量可控、历史语义稳定、重算频率低、实时和跑批的时效差异不大。否则不要轻易把重算的希望全押在流队列上。2.3 流批统一的新思路一套代码两式运行近几年的演进方向是用同一个引擎和同一套SQL同时支持流和批两种执行模式。Spark Structured Streaming是这么设计的Flink更是从1.x开始就把Table API和SQL作为统一抽象。你在Flink上写一套INSERT INTO语句把source改成Kafka就是流式写入改成Hive、Hudi或Iceberg文件就是批量写入。逻辑一样只是执行模式不同。这个思路在工程上被叫做“流批一体”。不是取消批和流的区别而是在开发、调度、计算、元数据层面共用一套体系。核心价值有三个第一代码开发量减半第二口径天然统一因为逻辑是同一段表达第三可运维性提升不用同时维护两套调度和两套监控。真正把流批一体落到生产环境光靠Flink或Spark本身还不够还需要一个能同时被流、批访问的数据存储层。Hudi、Iceberg、Paimon这类数据湖表格式就是在这种背景下被大规模采用的。把实时写入的结果先以增量形式落到湖表批处理在湖表上做批量读取两边看到的都是同一张表既保证实时性又能做历史切片查询。所以在实际选型的时候我并不觉得Lambda和Kappa是二选一的对立物。更多时候大家是在两者的框架之上找一条“用统一代码做两态运行”的路径。下面的内容重点讲我判断选型时实际会用的方法。3. 选型判断别急着定架构先回答这七个问题我见过太多团队在项目启动会上直接争Lambda还是Kappa吵两个小时没有结论。后来我自己总结了一个做法先别谈架构把业务的技术约束列清楚用一组问题来做筛选。这比看任何架构大会的PPT都管用。3.1 七个必须回答的问题这七个问题是我和团队在选型初期一定会过一遍的实时性要求到底是多少是秒级、分钟级还是小时级就行数据量级每天多大高峰QPS大概什么水平消息队列能撑住多久的保留计算逻辑复杂度高不高有没有大窗口、多流Join、长周期状态这类对状态存储压力很大的算子历史重算的频率高不高如果业务口径变化多久会触发一次全量或近全量重算实时和离线指标是不是必须同口径、同输出业务侧对两边数字不一致的容忍度有多大团队的人力和技术栈是什么情况是Flink和Spark都会还是只熟一套现在已有的链路复杂度如何如果是老旧Lambda能不能接受逐步迁移的过渡期这七条看下来大部分团队其实已经有倾向了。我直接给几个典型画像大家可以对号入座。3.2 什么样的场景适合直接上Kappa风格如果你的数据量每天在几百GB以内Kafka保留两三天的数据成本可接受业务重算需求极低团队又主要熟流引擎那可以走Kappa。典型业务比如实时推荐特征、在线学习样本拼接这类场景对精确历史回溯要求不高但非常在意延迟和代码简洁性。还有一个常见误判以为用Kafka存了全量数据就是Kappa。存储永远不是问题核心真正的关键是你有没有能力在有限时间内把全量历史重新消费并计算完。判断标准很简单把一年数据量除以重放时间上限再除以消费并行度估算出的吞吐是否在集群承受范围内。算完之后很多团队会主动放弃Kappa。给一个粗略的估算例子假设一年数据量约100TB要求重放在12小时内完成Kafka单partition消费吞吐按20MB/s估算一天86400秒理论吞吐约1.3GB/s换算约1.7TB一小时跑12小时约20TB。这离100TB差得远。要达到目标需要大概5倍以上的吞吐资源。如果预算不支持Kappa就直接排除。3.3 什么时候应该保留Lambda思想反过来如果你有这些特征建议保留Lambda思想甚至在它上面做改良业务方明确要求离线报表和实时大屏数字完全对齐一点都不能差已经有复杂的存量Hive数仓体系SQL逻辑几百条短期不可能翻写有较大的状态规模或者按用户维度的长周期聚合流引擎状态管理有压力这些情况下比较务实的做法不是硬拆掉Lambda而是用流批一体的方式降低双链路的成本。比如把实时和离线用同一套Flink SQL表达离线调度的结果仍然走批链路出报表实时链路只输出亚秒级指标两边用同一张数据湖表作为底表靠一个“对账任务”定时比对实时结果和批结果的差异。L这种“双跑但共用口径”的改良形态是我在团队里推荐最多的方案。架构不是信仰成本收益比才是核心。4. 落地实操从Lambda平滑迁移到流批统一4.1 先定底座消息总线、计算引擎、湖表三件套在真正开始迁移之前先把三件事定下来消息总线用Kafka还是Pulsar计算引擎用Flink还是Spark Structured Streaming表格式用Hudi、Iceberg还是Paimon。我的建议是团队熟悉哪个用哪个没有绝对最优。我个人的偏好是Kafka Flink Hudi这套组合的生产资料最多社区排障案例也够丰富。Kafka失序靠key分区基本能解决Flink的checkpoint机制成熟Hudi在流式写入的upsert和增量读取上都算稳定。Iceberg的流批一体能力也很好但当时实时upsert和文件合并不如Hudi成熟。Paimon则适合希望把流式主键表和批式数仓进一步打通的团队它更年轻踩坑要靠社区活跃度顶住。这套底座确定之后目录规划也很有讲究。建议直接以库表隔离的方式管理层次比如ods层实时源数据直接落Hudi表保持与Kafka消息基本同步dwd层流式清洗、维度补全后的明细表dws层汇总指标表流批共用表名和字段统一露出既是流批统一的物理基础也是以后出问题好定位的关键。4.2 核心实现一套Flink SQL吃下流和批用Flink SQL写流批一体关键是同一套逻辑不要写两遍。举一个很常见的例子订单事实表的加工逻辑既要实时更新当天订单汇总又要支持离线批量回刷历史。流式写入Hudi的核心实现类似这样CREATE TABLE order_dws_daily ( order_date STRING, order_status STRING, order_cnt BIGINT, gmv_amount DECIMAL(20, 2), PRIMARY KEY (order_date, order_status) NOT ENFORCED ) WITH ( connector hudi, path hdfs:///warehouse/dws/order_dws_daily, table.type MERGE_ON_READ, hoodie.table.name order_dws_daily ); INSERT INTO order_dws_daily SELECT DATE_FORMAT(order_time, yyyy-MM-dd), order_status, COUNT(*) AS order_cnt, SUM(order_amount) AS gmv_amount FROM kafka_order_dwd GROUP BY DATE_FORMAT(order_time, yyyy-MM-dd), order_status;这个写入任务跑起来之后Hudi表里的数据是带有实时增量的。批处理侧要读取同一张表只需要把执行环境切换成批模式Configuration config new Configuration(); config.setString(execution.runtime-mode, BATCH); StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(config);或者如果在Flink SQL Client环境执行SET execution.runtime-mode BATCH; SELECT order_date, order_status, SUM(order_cnt) AS total_cnt, SUM(gmv_amount) AS total_gmv FROM order_dws_daily GROUP BY order_date, order_status;注意批模式读到的数据和流模式写入的数据源是同一个表不需要再做一次INSERT INTO到另外的Hive表。这套逻辑里实时口径和离线口径天然一致因为它们用的就是同一段加工表达式。4.3 生产链路从开发到上线要盯的细节很多团队在写SQL那一步很顺利反而在上线之后被细节拖垮。我最常让团队检查的有以下几件事第一Hudi表类型的选择。COWCopy On Write适合读多写少查询性能好MORMerge On Read适合写多读少写入吞吐高但读的时候可能要合并log文件查询延迟略高。流式upsert场景我更推荐MOR因为它对高频小文件的I/O压力更小。但如果下游主要是频繁的ad-hoc查询COW可能反而更好。这个没有定论需要实测读写链路再定。第二checkpoint和状态配置。建议先把Flink checkpoint interval压到1到3分钟incremental checkpoint打开。状态比较大时开启本地恢复让状态从本地磁盘恢复而不是全量从远端拉。还有处理Kafka消息的key要想清楚用订单ID还是用户ID决定了下游去重和Join的正确性。第三小文件和compaction策略。Hudi的同步compaction开开小文件阈值调小一点。我一般把最大文件大小设为128MB到256MB把提交清理策略打开。否则跑一周之后文件数量爆炸查询性能直线下降。第四监控和告警。不要只看作业状态。Flink的Kafka消费延迟、checkpoint失败次数、Hudi的写入事务失败率这三个指标一个都不能少。我自己会额外加一个“流入vs落表”的行数校验拿Kafka topic昨天消费总条数和Hudi表昨天新增条数做比对差值超过0.1%就告警。这个告警不复杂但能挡掉绝大多数链路异常。4.4 迁移节奏小步走双跑过渡完整的Lambda平滑迁移我强烈建议按以下节奏来第一步选定一个核心链路做试点最好是逻辑清晰、口径争议少的。不要一上来就把全量数仓翻写。把生成的Hudi表和已有Hive表做并行写入两边数据跑一个月每天都比对。第二步业务侧确认两边数据能对上了再分批切换读路径。可以先让报表查询走Hudi再让下游血缘依赖逐步切过去。旧的Hive表任务暂时保留避免上游血缘断掉导致的连锁反应。第三步新的流批链路稳定运行一段时间后把旧的批处理链路下线。下线动作必须提前半个月通知全部依赖方并且留出至少一个完整的月结周期作为观察期。别信“这个任务没人用了”这种话一定以实际血缘为准。这套节奏走下来基本能把从Lambda迁移到流批统一的风险压到最低。5. 躲坑实录这些坑我替你踩过了5.1 高频问题与排查思路这里直接上排查速查表都是我或者团队同事在实际项目中碰到的、有一定共性的坑。问题现象可能原因排查路径解决建议实时和离线数字对不上流批逻辑不一致或去重口径不同逐层对比dws层明细找到第一条不一致记录统一加工SQL把去重逻辑下沉到同一段表达式写入Hudi后查询越来越慢小文件过多或未开启同步compaction检查文件数量、Hudi timeline中compaction完成情况开启同步compaction调小commit间隔让文件更快合并批量重放历史数据时状态不一致上游Schema有历史变更对比消息schema变更记录与重放起始位点定义统一的schema演进规则重放前做好字段映射Kafka消费积压作业吞吐上不去单partition并行度受限或key分布不均查看各partition的lag对比作业并发度按业务key重新分区提高consumer并行度checkpoint频繁失败状态过大或写入外部系统超时查看checkpoint失败时的异常堆栈和背压情况开incremental checkpoint检查Sink端连接池和超时设置批处理读取Hudi读不到最新数据文件尚未commit或时间线可见性配置不对检查Hudi commit时间线和最新数据文件时间戳开启time travel或调整读取模式配置如begin.instant5.2 一次真实的踩坑复盘有一个比较典型的项目是做用户行为分析平台。最初是标准Lambda离线部分用Spark每天晚上跑全量用户路径分析实时部分用Flink接Kafka输出最近一小时的热门事件。两套代码是不同同事维护的口径差异特别大——离线定义的“活跃用户”是当天有登录或有浏览行为的用户实时那边只算了有点击行为的用户导致实时大屏永远比离线报表低一截。后来我们做了第一次统一把实时和离线逻辑都收敛到Flink上写成了同一个“用户行为标准化”SQL先对点击、曝光、登录三类时间做union再按用户维度聚合。统一之后两边数字终于能对齐了。但这还没完真正的坑在后面。新逻辑上线后实时链路没问题但离线批处理一到月初做整月回刷就跑得很慢。后来才发现问题出在Hudi表采用了MORlog文件积压太多批处理Scan时把log全量合并了一遍。当时只想着写入快忽略了读侧合并的代价。排查思路倒是不复杂先看Hudi表是COW还是MOR然后看小文件数量、log文件的大小和数量、compaction的触发频率。我们最后把compaction策略改成“每10个commit触发一次”的方式又调整了文件大小上限批处理读取性能才算恢复正常。这个案例最大的教训是选型不是看单点性能而是要看整个链路的读、写、查、重放四件事是否协调。你只优化写读就会拖后腿只优化读写入就卡住。流批统一方案不是加一个Flink就完事的它是一个整体工程设计。5.3 我的几个实战心得最后分享几个在多家公司都验证过的实战心得篇幅有限挑最实用的几条。关于Kappa的历史重放如果不得不做重放强烈建议把重放任务写成一个独立版本不做任何幂等之外的依赖。比如重放时临时关掉对推进下游的发送只在落表阶段做upsert。这样能避免重放时把一堆脏数据推到下游在线存储里。关于对账任务这是所有改造里性价比最高的一个东西。用Shell每15分钟调度一次拿Kafka累计消费条数和Hudi新增条数做对比差异超过阈值就发钉钉告警。一开始可能会觉得烦人但一年下来你会发现它挡掉的故障比你想的多得多。关于团队技术栈如果团队对Flink的Table API和SQL不太熟不要强行上全套Flink SQL流水线。可以先从消息接入、简单清洗开始一边熟悉一边扩大范围。流批统一的改造本质上是技术栈的演进不是一次性替换。这些话听起来不像是架构理论的干货但真正被生产环境折磨过的人应该知道这每条背后都是实打实的时间和事故换来的。
返回列表