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

资讯详情

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

从批量到实时:大数据处理架构演进与选型实战

从批量到实时:大数据处理架构演进与选型实战 做大数据处理这些年我面试过不少人也被问过无数次同一个问题批处理和实时处理到底怎么选说实话这个问题的答案一直在变。十年前能跑批就能解决绝大部分需求五年前谁家没一套实时链路出门都不好意思跟人聊数仓发展到现在再纠结批量还是实时已经有点像在问“米饭和面条哪个更好吃”——不是二选一而是看什么时候吃、配什么菜。这篇是这个系列的第二十三篇正好借“从批量到实时”这个主题把批处理的底层逻辑、实时链路的组件构成、架构的演进路线以及我在实操中踩过的坑、验证过的方案一次性梳理清楚。内容不追求大而全重点是给正在做技术选型、或者正在被存量离线任务折磨的同学一些可以直接参考的思路。1. 批处理大数据计算的起点与长期主力1.1 批处理的核心思想攒一批算一次批处理的核心思想和食堂做饭一个逻辑不会来一个客人就开一次火而是到了饭点把一批食材一次性加工完统一出餐。这么做的好处是效率高、成本低坏处是不能随到随吃错过饭点就得饿着等下一批。大数据里的批处理也是这套逻辑。以经典的Hadoop MapReduce为例数据先落到HDFS上任务启动后Map阶段读入数据、拆分计算Reduce阶段汇总结果整个过程以“作业”为单位运行作业跑完结果落盘。后来的Spark Batch也是同一套思想只不过把中间结果尽量放内存靠RDD和DataFrame把跑批的速度提升了一个量级但本质依然是“先攒数据再统一算”。这种“先攒后算”的模式优点非常突出。计算过程容易重试和重跑数据出错了可以回滚到前一天重新执行系统对资源的要求也不苛刻半夜低峰期把集群占满也没人管审计起来简单每个任务跑完都有日志和结果表。所以T1报表、离线数仓、离线训练样本生成这类场景直到今天依然重度依赖批处理。有意思的是这种“批量”思想不只存在于大数据体系里。日常办公里的批量重命名、批量改图片尺寸、批量下载文档用脚本批量处理CAD图纸、批量抠图、批量替换PPT素材底层思路都是共通的把重复劳动合并成一次批量操作用机器的吞吐换人的时间。理解了这一点再回头看大数据批量计算就没那么玄了。1.2 批处理的技术栈与适用边界批处理的技术栈这么多年下来基本稳定在四个层面存储层HDFS、对象存储S3/OSS、数据湖格式Iceberg/Hudi/Paimon的离线写入路径计算层MapReduce基本退到幕后、HiveSQL化跑批的老将、Spark批处理目前离线计算的主力调度层Crontab小规模场景、Airflow、DolphinScheduler重调度、带依赖管理、带告警服务层数仓建模后的OLAP引擎比如Doris、ClickHouse、StarRocks供查询分析使用这套组合能成为过去十几年的主流核心原因在于大部分业务决策并不需要秒级数据。运营看昨天的GMV、财务看上个月的成本、模型训练用过去一周的样本这些场景天然允许延迟。批处理刚好以最低的成本、最高的稳定性满足了这部分需求。但批处理的代价也摆在明面上延迟高通常以小时为单位数据新鲜度差业务没法做实时干预而且为了追求吞吐链路往往设计成“大任务、重计算”对突发情况反应迟钝。这也是为什么后来实时处理会崛起——不是批处理不行了而是业务对“当下”的要求变高了。提示批处理不会消失只会被拆得更小。很多号称实时的系统本质是把“一天一批”改成“五分钟一批”这在工程上叫微批后面会专门聊到。2. 实时处理从“事后”到“当下”的范式跃迁2.1 实时需求从哪里来一句话总结凡是数据晚到几分钟就会损失价值或金钱的场景就有实时处理的诉求。最典型的是实时特征服务。做推荐、风控、广告的同学应该深有体会——用户刚点了一个商品模型要立刻根据这个动作更新推荐结果如果特征链路还是T1模型看到的永远是昨天的用户转化率自然上不去。特征值越新鲜模型效果越好这就是实时的直接业务价值。类似的场景很常见股票的实时行情和K线数据量化策略要基于最新价格做决策运营大屏要实时刷新订单量和在线人数老板打开屏幕就要看到当前状态安全风控要根据用户行为实时判断是否盗号、刷单数据库同步要从“每天拉一次全量”变成“每次变更实时同步”。这些需求都指向同一个方向——把处理延迟压缩到秒级甚至毫秒级。这里顺便澄清一个概念边界。大家搜实时相关问题时会看到很多不同语境下的“实时”操作系统内核的实时补丁、IDE的实时调试、杀软实时保护、实时语法校验、实时抓包。这些“实时”和数据处理的“实时”完全是两回事。前者侧重系统对事件响应时间的确定性保障后者侧重数据从产生到可被消费的时间间隔。做大数据实时处理关注的是延迟、吞吐、准确性和可用性别被概念混淆带偏了。2.2 实时处理的核心组件与链路实时处理不是某一个组件能搞定的而是完整的一条链路。从生产到消费大致分成四个环节采集把数据从业务系统实时搬出来。日志类的用Flume、Logstash、Filebeat数据库变更类的用CDCChange Data Capture工具比如Debezium、Flink CDC。CDC的原理是监听数据库的binlog/redo log把insert、update、delete操作实时转发出来。传输采集到的数据需要一个中转站最常用的是Kafka或Pulsar。消息队列的作用不只是临时存储更重要的是削峰填谷和故障隔离——上游流量再猛下游处理不过来时可以慢慢消费不至于直接压垮服务。计算实时计算引擎负责对数据做清洗、关联、聚合、窗口计算。目前的事实标准是Flink其次还有Spark Streaming本质是微批、Storm老牌选手已明显边缘化。Flink能胜出核心在于把“精确一次语义”“事件时间处理”“状态管理”这些实时计算最头疼的问题做成了内建能力。服务计算完的结果要能支撑查询或决策。实时数仓一般配Doris、ClickHouse这类OLAP特征类场景配Redis或专门的特征存储指标类场景配Prometheus或自研监控系统。把这四个环节串起来就是一条标准的实时数据处理链路业务数据变化 - CDC/日志采集 - Kafka - Flink实时计算 - 结果写入OLAP/缓存/消息队列 - 业务查询或触发动作。下面用一张表把批量和实时做一个直观对比做方案选型的时候可以直接对照参考。对比维度批量处理实时处理延迟分钟到小时级秒级到毫秒级处理模式攒一批再算来一条算一条或微批吞吐高适合全量大扫描相对受限需要横向扩展容错重跑容易重跑即可困难需要状态恢复和幂等设计资源成本低可利用低峰资源高链路常驻运行典型场景离线报表、模型训练特征服务、实时告警、实时大屏3. 架构演进Lambda、Kappa与流批一体3.1 Lambda架构两条腿走路的岁月实时需求刚冒头的时候业界的做法很直接在原有批处理旁边再盖一条实时处理链路。批处理继续算T1的精确结果实时链路先出一个“快但可能不够准”的结果最后在一层对外服务里把两条链路的结果合并起来。这个模式后来被命名为Lambda架构。Lambda架构的好处是落地快、风险低坏处也相当明显同一套业务逻辑要写两遍——一遍用批处理引擎实现一遍用流处理引擎实现。两套代码、两套调度、两套资源踩得最多的坑是“口径不一致”。批处理算出来GMV是100万实时链路算出来98万业务方问差在哪往往要排查一下午的时间范围、去重逻辑、时区转换问题。维护成本极高这也是那个阶段最让人头疼的地方。3.2 Kappa架构用日志回放换掉重算为了化解双链路带来的重复开发和口径问题有人提出了Kappa架构。思路很直接既然实时流处理引擎本身就是在不断消费消息队列里的数据那为什么不让它在历史日志上直接回放把消息队列的保留时间拉长新任务上线时从最早的消息开始重新消费一遍相当于用流处理引擎完成批处理的事最终只保留一套代码。Kappa架构在逻辑上非常优雅但落到工程上同样有限制如果历史数据量极大、需要回溯三个月甚至一年的日志消息队列的消费速度可能远慢于批处理引擎直接扫描HDFS而且长期占用消息队列的存储成本也不低。所以Kappa并没有完全取代Lambda反而让业界开始思考一个更务实的问题——能不能在存储和API层面统一批流让开发人员只写一套逻辑系统自动决定什么时候按批处理、什么时候按流处理。3.3 流批一体现在的可行解流批一体简单说就是“一套存储、一套API、两种执行模式”。存储层面以Iceberg、Hudi、Paimon为代表的数据湖格式可以把实时写入的数据按事务方式落成列式文件批任务可以直接扫描同一份数据。计算层面Flink从很早就提出“流批一体”同一段SQL可以在流模式和批模式下运行Spark也在不断加强流批能力。实际落地中我见过最多的形态是业务库 - Flink CDC - Kafka - Flink - 实时结果表/数据湖 - 实时OLAP供查询同时批任务直接读取数据湖做全量分析。这个形态的好处是实时链路和离线链路共享底层同一份数据口径天然的拉齐了一大半剩下的就是窗口边界和时区规范问题。选数据湖格式的时候我的经验是重点关注三件事一是事务能力和并发写控制实时写入时能否保证读到的数据是快照一致的二是小文件治理能力实时任务写入频率高小文件累积非常快Paimon和Hudi都内置了compaction机制来缓解三是生态兼容性能否被Spark、Flink、Trino同时读写。如果陷进某个引擎的私有格式里后面想换引擎就非常痛苦。4. 落地实践三个可以直接上手的小场景4.1 场景一实时特征服务与模型实时评分先聊一个和我日常关系最密切的场景实时特征服务。典型链路是用户行为日志打进KafkaFlink消费后实时计算累计点击、最近N分钟浏览时长等特征写入Redis或特征存储在线推理服务每次收到请求时从特征存储实时拉取用户特征喂给模型比如逻辑回归打分。这个链路解决的痛点是模型参数可以离线训练好但输入特征必须在线更新否则用户刚发生的行为不会被计入本次决策。用Python模拟一下这个思路核心逻辑并不复杂重点是“消费一批事件持续更新特征状态”import time import redis r redis.Redis(hostlocalhost, port6379) def process_event(event): # event: {user_id: u123, action: click, ts: 1700000000} uid event[user_id] # 累加总点击、5分钟内点击等指标 r.hincrby(fuser:{uid}, total_clicks, 1) r.hincrby(fuser:{uid}, click_5min, 1) while True: # 生产环境这里从 Kafka 拉数据这里是模拟 event fetch_from_kafka() if event: process_event(event) else: time.sleep(0.1)生产环境当然不会写得这么简单但核心思想一致把特征从离线的预计算变成在线状态的持续更新。做这类系统最容易踩的坑是特征更新的一致性——同一用户的两个并发事件同时写入要确保用事务或Redis的原子命令比如hincrby避免状态错乱。这类服务对延迟极其敏感特征写入路径上的任何阻塞都会直接影响推理时效所以能放内存和缓存的状态尽量别落在关系型数据库里。4.2 场景二Python实时可视化刷新另一个经典落地是实时大屏。现在很多团队习惯用Python做数据分析和快速原型可视化也常用。技术选型上有三条路最粗暴的是前端定时轮询每秒刷一次接口好一点的是用WebSocket或SSEServer-Sent Events服务端有数据变更就主动推送再进一步就是用Streamlit这类框架把实时更新的逻辑直接写进脚本里。我来个简化的Python示例演示“数据不断进来、页面自动刷新”的核心思路import streamlit as st import pandas as pd import random import time placeholder st.empty() while True: new_row {value: random.randint(1, 100), time: pd.Timestamp.now()} # 真实场景这里从 Kafka/Redis 读取实时数据 st.session_state.setdefault(data, []).append(new_row) df pd.DataFrame(st.session_state[data][-100:]) placeholder.line_chart(df.set_index(time)) time.sleep(1)这个写法适合做快速原型和内部看板。生产级大屏一般还是会用前端ECharts或Grafana结合WebSocket来做性能和交互都会好很多。如果刚接触实时可视化我建议先用Streamlit或Grafana这类现成工具跑通链路再考虑自研前端别一上来就上全家桶。4.3 场景三存量批量作业的实时化改造最后聊一个更普遍的需求把现有的批量任务改造成实时链路。最典型的例子就是数据库同步。过去做数仓大部分是从Oracle、MySQL每天凌晨拉一次全量第二天早上出报表现在业务要求看到今天中午的数据就得把“每日全量同步”改成“实时增量同步”。工具选型上优先考虑支持CDC的方案如果用的是云数据库直接看云厂商提供的数据传输服务自建数据库就用Debezium加Kafka这是目前社区里最通用的组合生态相对完整下游可以接Flink也可以接专门的数据同步服务。实时化改造有个务实的做法不要一步到位。我的建议是分三步走。第一步同步层先实时化用CDC把数据实时落到Kafka或数据湖这一步成本低、见效快。第二步清洗和关联层从每天定时JOB改成Flink流式作业期间和旧链路双跑逐日对比结果差异。第三步等新链路稳定后再切流量、下线旧任务。每一步都有可回退的余地不至于一把梭把整个数仓搞乱。注意实时化改造最大的风险不是技术而是“需求是不是真的需要实时”。如果业务方只是觉得“实时听起来高级”那改完之后大概率会陷入无穷无尽的口径对齐和资源投入。上实时之前先问清楚用户到底要秒级还是分钟级很多场景其实“五分钟微批”就够了。5. 常见问题与排查技巧实录5.1 实时链路里最坑的五个问题写实时作业最常见的坑我整理成了下面这张表问题表现常见解决思路乱序与迟到数据结果不对、数值跳变设置合理的Watermark策略用事件时间而非处理时间重复消费同样的数据被算了两遍下游幂等写或用Flink Checkpoint配合端到端一致性背压上游生产速度超过下游处理能力定位瓶颈算子增加并行度或优化状态访问窗口漂移/时区问题结果和离线表对不上统一时区规范明确窗口边界建议用UTC计算再本地化展示实时写湖小文件文件数暴涨、查询变慢开启compaction或对写入做微批缓冲比如每30秒提交一次乱序和迟到数据是最容易忽略的。实时链路里网络抖动、服务重试都可能导致事件先后顺序颠倒如果直接按接收顺序累加数值就会忽高忽低。正确做法是用事件自带的时间戳比如用户点击时间做计算配合Watermark判断“迟到的消息还等不等”。重复消费的问题同样隐蔽。Kafka的offset机制在正常情况下能做到至少一次但程序重启、网络异常时很难保证不多算一遍。如果下游是计数类指标建议做幂等写入用唯一键或版本号让重复数据在写入层被自然过滤。这个经验是踩了好几次坑才得来的。5.2 批量转实时的过渡期避坑清单如果团队正处在批量转实时的过渡阶段有几个点强烈建议重视双跑验证新老两条链路并行跑至少对比一周。对比维度包括总量、明细、高峰时段趋势别只盯最终汇总数。灰度切流按先内部后外部、先低敏后高敏的顺序切换数据消费方切流期间保留回滚开关。监控指标实时链路的监控比离线更依赖细致指标体系。延迟、积压量、消费速率、Checkpoint失败次数、结果偏差率最好都做成可视化看板。数据回看实时链路正确性不能只看当前值定期用离线全量数据做“回看对账”确保实时结果和离线结果统计口径一致。见过太多实时数仓上线三个月实时和离线对不上最后排查发现是窗口边界定义不一致。再说一个容易被轻视的点实时任务的资源规划和成本预估。实时链路是7x24小时运行的不像跑批可以半夜集中用资源。方案设计阶段就要评估清楚Kafka的存储时长、Flink的常驻资源、OLAP引擎的并发度否则上线一个月的账单可能让你重新思考什么叫“实时”。6. 我的一点真实体会做到今天这个阶段我的核心体会是批量不会消失实时也不是银弹。大部分公司最终的形态是“批量 实时”并存关键是让它们收敛到同一套存储、同一套口径上。实时解决新鲜度问题批量解决准确性和成本问题两者不是替代关系而是分工关系。用我自己的话说数据处理的演进与其叫“从批量到实时”不如说“从单一时效到分层时效”。T1能解决的场景继续用批量分钟级能接受的用微批只有真正需要秒级响应的场景才值得上全套实时链路。决定上实时之前先回到那个最经典的问题数据晚到十分钟业务会不会有实际损失如果不会趁早别折腾。如果说有什么经验是最后想分享的那大概是实时化改造这件事最容易犯的错误不是技术选型错了而是低估了运维复杂度、高估了业务收益。以最小的代价满足需求留出足够的扩展空间比任何一步到位的方案都靠谱。
返回列表