早十年做数据开发的同行,应该都经历过这样的日常:白天业务系统不停往里写数据,凌晨调度任务开始跑批,天亮之前把报表算完,第二天业务方看到的永远是昨天的数据。这套“批量处理”的打法统治了数据圈很多年,稳定、可控、出错能重跑。但最近几年风向确实变了——越来越多团队在从“T+1给答案”转向“秒级出决策”,大数据处理从批量到实时的技术演进,成了绕不开的话题。今天这篇不想讲空概念,就当一个做了多年数据的老兵跟同行复盘:批处理为什么能活这么久,实时到底在解决哪些实际问题,以及真正做演进时你会踩到哪些坑。
顺便说一句,我经常看到有人在网上搜“批量修改文件名”“批量重命名”这种脚本,说明“批量”这个思路在日常工具领域还是刚需;但放到大数据处理里,趋势已经明确指向“实时”。这个反差特别有意思——不是批量没用了,而是数据量一旦大到某种程度、业务决策一旦要跟秒级挂钩,批量就撑不住了。
1. 批处理这套“黄金打法”,是怎么统治数据圈的
1.1 批处理的思想根源:把“算力”攒起来一次性用
在聊实时之前,得先把批处理讲透。批处理的核心思想特别朴素:数据先落库,攒到一定时间窗口或数据量级,然后一次性读出来全量计算。计算过程中输入数据不变化,结果稳定、可预期。
批处理不等于落后。它最大的优势就是确定性。输入是固定的静态数据集,第一遍算和第十遍算结果完全一致;如果某一层算错了,把输出删掉重新跑一次就行,不会污染上游数据。这个特性在早期链路不稳定、数据质量没有保障的环境下,简直比什么都重要。反过来看,如果数据像自来水一样时刻流动,你想“重跑”,还得先弄清楚水流到哪儿了、哪些数据已经算进去了,复杂度完全不是一个量级。
我见过不止一个团队,搞所谓“实时化”之后,第一件事就是停掉原有的批处理任务,结果一遇数据问题就抓瞎——实时链路里没有“重放按钮”,或者说重放的代价高到没人敢轻易按。所以这篇文章首先想传达的一个态度是:不要把批处理当成敌人,它是你演进路上的安全垫。
1.2 从MapReduce到Hive:离线数仓的成熟路径
批处理在工程上的成熟,基本就是Hadoop生态的发展史。最早的MapReduce是纯Java编码,业务逻辑再简单都要写一堆map、reduce函数,迭代效率很低。Hive出现后,把SQL翻译成MapReduce任务,才算真正把数仓普及开。再后来Spark用内存计算把中间结果留在内存里,原本要跑一整晚的T+1批处理,缩短到几个小时。
这个阶段其实已经出现“批量到实时”的萌芽:不是业务变实时了,而是计算变快了,让原本隔天才能出的结果,能在数小时内看到。但本质上,它仍然属于批处理——数据是切片式的、任务是定时触发的,只是切得更细、跑得更快。
这给了一个很重要的启示:批量到实时的演进,往往不是一步到位,而是先把批处理做快,快到一个阈值之后,再考虑要不要换引擎。
1.3 批处理的三种典型任务形态
批处理任务在工程上大概有三种常见形态:
- 定时全量:针对数据量小、变化不频繁的维度表,每天或每周整体重建一次。
- 增量合并:按日期分区或按binlog偏移量,把新增数据合入数仓目标表,日常任务绝大多数都是这种。
- 重跑修复:因为批处理结果可预期、可重复,出了数据问题可以只删掉错误分区,定向重算修复。
三种形态都依赖一个前提:数据是“有边界”的。只有当一批数据彻底结束,你才能开始计算,而批次周期决定了延迟上限。这就是批处理不能做实时最根本的原因,不是算得不够快,而是它的执行模型里天然有个“等待数据结束”的步骤。
2. “实时”到底在解决什么:拆开业务需求层层看
2.1 真实业务里,谁在真正喊“实时”
这些年喊着要实时的业务,我大致归成几类,代码里的复杂度和成本完全不一样。
- 风控类需求最刚性:盗刷、欺诈、异常登录,每一秒都在发生,你不可能跟业务说“明天汇总完看看有没有问题”。必须交易发生时毫秒级判定,这直接决定了数据库同步、特征计算、模型推理要全链路实时化。
- 实时大屏内容是展示刚需:数字孪生、工业监控、运营驾驶舱,核心是“此刻”的状态,不是十分钟前的快照。这类场景容忍少量延迟,但追求持续刷新。
- 实时特征服务是机器学习在线化的标配:模型离线用T+1样本训练,但推理时必须喂入实时特征,比如当前设备指纹、最近一小时行为序列。这里就会出现热搜里常提到的“实时特征服务”,本质上就是批处理和实时处理交界的地方。
- 金融行情是天花板级别需求:实时K线、盘中预警、行情推送,一秒都等不了,对全链路延迟要求极高,普通大数据团队不一定碰得到。
2.2 数据新鲜度的四个层级
很多人一提“实时”,默认就是毫秒级。其实数据新鲜度至少有四个层级,每个层级的技术栈和成本差异巨大:
| 层级 | 数据新鲜度 | 代表方案 | 成本 |
|---|---|---|---|
| T+1离线 | 24小时以上 | Hive/Spark离线调度 | 低 |
| 准实时 | 10分钟~1小时 | 短周期调度+增量同步 | 中低 |
| 近实时 | 1~5分钟 | 微批流处理 | 中 |
| 实时 | 秒级~毫秒级 | Flink/事件驱动/CEP | 高 |
这里要特别说一句:把“准实时”做好,对大多数业务已经足够。很多报表和Dashboard做不到秒级的原因不是没有好引擎,而是需求本身不需要。过度设计是数据团队最容易犯的毛病,后文单独讲。
2.3 实时的共同本质:从“事后复盘”变成“事中决策”
把上面几类需求放一起看,会发现它们有一个共同点:把数据从“事后复盘”变成“事中决策”。批处理解决的是“昨天发生了什么”,实时解决的是“正在发生什么、要不要马上做点什么”。
理解了这一点,很多技术选型就不再纠结。比如一个场景如果业务根本不需要在事件发生的瞬间做决策,那你完全没必要引入Flink,用秒级或分钟级的定时任务就够了。真正的实时需求,往往都直接指向决策动作——拦截一笔交易、触发一个告警、推送一条消息、调整一次库存。没有“决策”这个动作的实时,多半是在自嗨。
3. 从T+1到T+0:架构演进的四个关键跳变
3.1 第一跳:从“全量重算”到“增量同步”
早期想把数据实时化,很多团队的第一反应是“把同步频率调高”。但全量同步在数据量上来之后必然扛不住,几千万行每天全量扫一遍,数据库压力大到业务方投诉。于是CDC(Change Data Capture,变更数据捕获)成了主流方案。
MySQL场景下,基本就是解析binlog。工具层面,我最常用的是Canal和Debezium。Canal在国内生态更熟,文档多,团队接手快;Debezium的好处是跟Kafka Connect集成天然,适合已经上了Kafka生态的团队。Oracle场景则要依赖LogMiner或专业同步工具。除了看支持哪种数据库,选型还要重点考察增量解析能力、断点续传、DDL兼容性,以及下游衔接是否顺畅——是直接进Kafka还是进Pulsar。这个决定影响很长时间,别只看演示demo跑通了就定。
从全量到增量,是批量到实时最关键的第一跳。不做这一步,后面所有实时计算都是空中楼阁。
3.2 第二跳:从“定时调度”到“事件驱动”
批量时代用Airflow、DolphinScheduler这类调度系统,核心模型是“到点触发”。但定时调度有一个天然缺陷:你永远没法精确知道上游数据什么时候准备好。任务定在凌晨2点,如果上游数据凌晨3点才到,这趟跑就白跑了;定太晚,又牺牲时效。
实时时代把触发模型换成了事件驱动——上游数据一有变化,立刻产生事件,下游消费到事件就触发计算。这个思维转变对很多数据开发来说是硬骨头,因为以前写惯了“每天几点跑”,现在要改成“数据产生了就去算”,还要考虑事件丢失、乱序、重复消费等一系列问题。
它的技术承载,就是消息队列加流处理框架。Kafka是最常见的选择,靠分区机制提供削峰和并行能力;Flink在流上做窗口计算和状态管理。演进到这里,批处理里的“调度依赖”渐渐不再是最核心的问题。
3.3 第三跳:从“单引擎”到“批流一体”
不少团队演进到一半会陷入一个尴尬:流处理用Flink写一套逻辑,离线批处理用Spark再写一套逻辑,两套代码指标口径经常对不上。这就是“批流一体”概念出现的直接原因——不再维护两套逻辑,而是用同一套引擎、同一套代码,同时处理有界数据和无界数据。
Flink从流处理起家,用DataSet API做有界流批处理;Spark则从微批出发,逐渐支持整批。选型看团队历史:如果已经有很重的Spark离线资产,硬切成Flink学习成本和迁移成本都很高;如果从零开始建实时能力,我更推荐直接用Flink,因为它的流处理语义和状态管理机制更贴近实时场景。
但批流一体绝不是银弹。它统一的是“逻辑表达”,并不能解决所有物理问题——比如离线要做大规模复杂Join,流式的状态压力不一定扛得住;有些历史SQL也不好平移。所以很多团队实际采用“流批两套引擎+公共口径层”的折中。
3.4 第四跳:从“离线数仓”到“实时特征服务”
最近两三年,我观察到最明显的一个趋势,是数据实时化开始往机器学习方向渗透,就是热搜里提到的“实时特征服务”。
传统模式下,特征是从离线数仓算好存起来的,模型训练和推理都用离线特征。但实时推理需要在线特征——用户刚点的这个按钮、刚发生的那次浏览,都必须马上进入特征计算并参与模型评分。一个典型的链路是:行为日志采集后经Kafka进Flink,Flink按实体ID做状态聚合,算出的实时特征写入Redis或Feature Store;业务侧调用特征接口时,把实时特征和离线画像特征拼在一起,喂给逻辑回归等模型。
这里出现的热搜词“逻辑回归实时评分主引擎scikit-learn 1.5.x实时推理”,其实就是这个形态的常见实现:离线用scikit-learn训练好模型,序列化后部署到在线服务,把实时特征拼装进去做预测。技术栈选型本身不复杂,真正的难点在特征口径的一致性、特征延迟的监控,还有数据回放补特征的能力。
4. Lambda、Kappa与混合架构:别再争了,得看场景
4.1 Lambda架构的经典形态与维护之痛
一聊实时架构,就绕不开Lambda。Lambda把链路拆成批层、速度层和服务层:批层用离线任务产出准确结果,速度层用流处理产出低延迟结果,服务层负责把两者合并后对外提供查询。
Lambda在逻辑上很清晰,但落地后痛点突出。最典型的就是同一套指标要在批和流里各写一遍,口径经常对不上。比如“今日成交额”,离线口径可能是“按支付成功时间统计”,实时口径可能变成“按订单创建时间统计”,两边差个尾数,业务方一问你就得去查半天。很多团队最终放弃Lambda,不是因为它不能工作,而是维护成本太高。
4.2 Kappa架构的理想与现实
Kappa是Lambda的简化版,主张只保留流处理一条链路,所有的历史数据重算都通过Kafka重放来完成。理想情况下,一套代码既处理实时的增量数据,也处理历史的重放数据,天然口径一致。
但落到工程上,Kappa有几个现实问题。第一,Kafka的消息保留时间有限,默认可能就是几天,做长时间的历史数据重放要么扩容,要么做数据落盘归档;第二,全量重放的耗时和成本远比离线任务高;第三,当数据本身就有质量问题,需要做历史修正时,流处理链路缺乏离线那种“删分区重跑”的便宜操作。所以纯Kappa在工单上都很好,真上线能坚持住的团队不多。
4.3 我实际推荐的折中:一份日志进两套引擎,用对账闭环兜底
我自己的经验是:不必非得在两个“派系”里站队。比较务实的折中做法,是让所有源头数据先进Kafka统一存一份,Flink实时算一份结果,离线任务(Spark或Flink Batch)从Kafka或数仓增量目录再算一份日级结果,两条链路共用同一份口径定义文件。
日常查询走实时链路,满足秒级场景;日跑结果用于对账和修正,发现差异后以离线口径为基准反推问题所在。这样既保证业务体验,又给数据质量留了一条退路。
| 对比维度 | Lambda | Kappa | 混合折中 |
|---|---|---|---|
| 逻辑复杂度 | 高,两套代码 | 低,一套代码 | 中,共享口径 |
| 口径一致性 | 难保证 | 天然一致 | 靠对账兜底 |
| 历史重放能力 | 强 | 依赖消息队列留存量 | 中 |
| 运维成本 | 高 | 中 | 中高 |
| 适合场景 | 重口径、重稳定 | 从零建设、小团队 | 有存量资产的演进型团队 |
我个人推荐有存量批处理资产的团队直接走第三条路,经验是“演进”而不是“革命”。
5. 落地实时链路时,最容易翻车的实操环节
5.1 采集端:日志的乱序与重复,比你想象的严重
很多第几节课跑通demo的实时链路,一上生产就暴雷,第一个坑通常在采集端。日志从业务服务发出后,经过网卡、代理、队列、消费者,顺序几乎不可能严格保持。不同服务产生的时间戳又不一致,机器时钟可能有偏差。
处理乱序的标准姿势是引入“事件时间”概念,配合Flink的水位线机制。但多少的乱序容忍度合理?这个参数要从业务实际压测,拍脑袋填一个3秒或5秒对某些场景就是灾难。此外,重复消费几乎无法避免,所以实时计算结果本身要尽量幂等,或者用状态去重。
5.2 Kafka分区键:不设计好,下游根本没顺序
Kafka在同一分区内是有序的,但分区之间不保证全局顺序。所以分区键的设计直接决定下游能否正确处理同一实体的数据。最常见的错误,是把所有数据都扔到同一个分区来保证全局有序,这样分区并行度完全浪费,吞吐量断崖式下降。
正确做法是按业务实体的ID分分区,比如按用户ID、订单ID取模。这样同一用户的事件一定落在同一分区、按顺序消费,不同用户之间天然并行。如果业务又要求全局有序,那要考虑是否真的需要全局有序——大多数场景的“全局有序”其实是伪需求。
5.3 Flink状态与Checkpoint:别把状态当黑盒
流处理里最难控的就是状态。很多人把状态当成一个看不见摸不着的黑盒,出了数据错乱不知道从哪查。Flink的状态后端是存在本地的,配合Checkpoint做快照,才能实现精确一次语义。实际操作中,Checkpoint间隔太短会造成大量IO压力,太长又让故障恢复时的回放时间变长。我建议起步按5到10秒间隔配置,运行后观察Checkpoint耗时和背压情况再调整。
并行度的设置也常被忽视。一个Flink作业的并行度不等于Kafka分区数,盲目调大并行度会让状态分布变碎,增加网络开销。比较稳的做法是让Kafka分区数与Flink并行度保持一致,然后再压测微调。
5.4 实时与离线对账:不做必翻车
实时链路跑顺以后,最容易丢掉的防线就是对账。我见过不止一次事故:实时大屏数字和财务离线报表差了几百万,业务半夜找过来,你却发现没有任何机制能自动发现差异,只能人工拉数核对。
对账机制不必做得特别复杂,先跑日级对账:每天凌晨用离线结果和实时结果做一次全量核对,覆盖关键指标;再跑分钟级监控:给核心指标设置波动阈值,一超过阈值就告警。一旦发现差异,排查链路一般是“源头日志是否丢→Kafka分区是否堆积→Flink窗口时间是否选错→幂等去重是否生效”。把这个清单事先写好,能省下大量救火时间。
6. 反直觉结论:不是所有大数据处理都必须“实时”
6.1 实时改造的真实成本
写到这里,我想唱个反调:不是所有业务都该上实时。实时改造的真实成本,很多人评估得太乐观。
第一是人力成本。流处理的开发和排错思路跟批处理差别极大,团队要学习窗口、状态、水位线、Checkpoint等一整套概念,磨合期至少一到两个月。第二是资源成本。Flink集群是要长期运行的,状态后端和Kafka都要占用大量存储;而批处理任务跑完资源就释放,夜晚低谷还能复用资源。第三是运维复杂度。实时链路故障恢复时间窗口苛刻,对监控、演练、值班的要求远超离线。第四也是最容易被忽略的,实时链路的故障恢复复杂度——一个凌晨的流任务堆积,就可能直接影响业务,而不是像离线任务那样延迟几小时才造成影响。
6.2 哪些场景用“近实时”就够了
如果业务数据量不是特别大,或者延迟需求在1到5分钟之间,用近实时方案往往性价比更高。比如分钟级调度配合增量同步,用Sparks微批或Flink的微批模式,既能简化状态管理,又能降低运维成本。
电商场景里,“昨日销售预测”完全不需要实时;真正要实时的,是“当前库存不足告警”“正在发生的恶意抢购”。把需求按决策紧迫性排个优先级,会发现真正需要秒级响应的,往往不到20%。剩下80%,用准实时甚至T+1就能解决。
6.3 按数据量、延迟需求、口径复杂度做决策
我常用的一个简易决策框架是这样的:
- 数据量小、延迟需求大于5分钟:不需要引入流计算引擎,用定时任务加增量同步即可。
- 数据量大、但可以接受分钟级:用微批或短周期调度,别为“伪实时”增加过多成本。
- 数据量大、要求秒级且强状态处理:才考虑完整的Flink流处理链路。
- 在线推理需求:实时特征服务是刚需,但特征计算可以先用规则做,再逐步上复杂模型。
这个框架不是精确公式,但它能帮团队避开一个常见误区:拿实时引擎去解决一个根本不需要实时的问题。
6.4 我的个人经验与建议
按我自己的实践体会,批量到实时的演进,最关键的是“需求驱动”而不是“技术驱动”。别因为看到别人都在搞Flink就觉得自己落后了,先问清楚业务方到底需要多新鲜的数据,需要哪个环节做决策,能接受多少延迟和成本。
如果确实要做,也建议走增量演进:先做增量同步和准实时链路,跑通后再平滑升级到秒级流处理。每次演进都要保留离线批处理这条安全通道,直到实时链路稳定运行几周。我自己踩过最深的坑,就是过早地把离线任务停掉,结果实时链路一故障,全公司报表直接瘫痪,最终花了整整一周才把离线口径重新对齐。从那以后,我立了一条规矩:实时链路要上线,离线任务必须继续跑一段时间,直到对账连续多日无差异,才允许逐步退坡。
这些经验不是教科书上教的,是拿真实故障换来的。希望对做数据开发、数据架构的同行有所帮助。