
早十年做数据开发的同行应该都经历过这样的日常白天业务系统不停往里写数据凌晨调度任务开始跑批天亮之前把报表算完第二天业务方看到的永远是昨天的数据。这套“批量处理”的打法统治了数据圈很多年稳定、可控、出错能重跑。但最近几年风向确实变了——越来越多团队在从“T1给答案”转向“秒级出决策”大数据处理从批量到实时的技术演进成了绕不开的话题。今天这篇不想讲空概念就当一个做了多年数据的老兵跟同行复盘批处理为什么能活这么久实时到底在解决哪些实际问题以及真正做演进时你会踩到哪些坑。顺便说一句我经常看到有人在网上搜“批量修改文件名”“批量重命名”这种脚本说明“批量”这个思路在日常工具领域还是刚需但放到大数据处理里趋势已经明确指向“实时”。这个反差特别有意思——不是批量没用了而是数据量一旦大到某种程度、业务决策一旦要跟秒级挂钩批量就撑不住了。1. 批处理这套“黄金打法”是怎么统治数据圈的1.1 批处理的思想根源把“算力”攒起来一次性用在聊实时之前得先把批处理讲透。批处理的核心思想特别朴素数据先落库攒到一定时间窗口或数据量级然后一次性读出来全量计算。计算过程中输入数据不变化结果稳定、可预期。批处理不等于落后。它最大的优势就是确定性。输入是固定的静态数据集第一遍算和第十遍算结果完全一致如果某一层算错了把输出删掉重新跑一次就行不会污染上游数据。这个特性在早期链路不稳定、数据质量没有保障的环境下简直比什么都重要。反过来看如果数据像自来水一样时刻流动你想“重跑”还得先弄清楚水流到哪儿了、哪些数据已经算进去了复杂度完全不是一个量级。我见过不止一个团队搞所谓“实时化”之后第一件事就是停掉原有的批处理任务结果一遇数据问题就抓瞎——实时链路里没有“重放按钮”或者说重放的代价高到没人敢轻易按。所以这篇文章首先想传达的一个态度是不要把批处理当成敌人它是你演进路上的安全垫。1.2 从MapReduce到Hive离线数仓的成熟路径批处理在工程上的成熟基本就是Hadoop生态的发展史。最早的MapReduce是纯Java编码业务逻辑再简单都要写一堆map、reduce函数迭代效率很低。Hive出现后把SQL翻译成MapReduce任务才算真正把数仓普及开。再后来Spark用内存计算把中间结果留在内存里原本要跑一整晚的T1批处理缩短到几个小时。这个阶段其实已经出现“批量到实时”的萌芽不是业务变实时了而是计算变快了让原本隔天才能出的结果能在数小时内看到。但本质上它仍然属于批处理——数据是切片式的、任务是定时触发的只是切得更细、跑得更快。这给了一个很重要的启示批量到实时的演进往往不是一步到位而是先把批处理做快快到一个阈值之后再考虑要不要换引擎。1.3 批处理的三种典型任务形态批处理任务在工程上大概有三种常见形态定时全量针对数据量小、变化不频繁的维度表每天或每周整体重建一次。增量合并按日期分区或按binlog偏移量把新增数据合入数仓目标表日常任务绝大多数都是这种。重跑修复因为批处理结果可预期、可重复出了数据问题可以只删掉错误分区定向重算修复。三种形态都依赖一个前提数据是“有边界”的。只有当一批数据彻底结束你才能开始计算而批次周期决定了延迟上限。这就是批处理不能做实时最根本的原因不是算得不够快而是它的执行模型里天然有个“等待数据结束”的步骤。2. “实时”到底在解决什么拆开业务需求层层看2.1 真实业务里谁在真正喊“实时”这些年喊着要实时的业务我大致归成几类代码里的复杂度和成本完全不一样。风控类需求最刚性盗刷、欺诈、异常登录每一秒都在发生你不可能跟业务说“明天汇总完看看有没有问题”。必须交易发生时毫秒级判定这直接决定了数据库同步、特征计算、模型推理要全链路实时化。实时大屏内容是展示刚需数字孪生、工业监控、运营驾驶舱核心是“此刻”的状态不是十分钟前的快照。这类场景容忍少量延迟但追求持续刷新。实时特征服务是机器学习在线化的标配模型离线用T1样本训练但推理时必须喂入实时特征比如当前设备指纹、最近一小时行为序列。这里就会出现热搜里常提到的“实时特征服务”本质上就是批处理和实时处理交界的地方。金融行情是天花板级别需求实时K线、盘中预警、行情推送一秒都等不了对全链路延迟要求极高普通大数据团队不一定碰得到。2.2 数据新鲜度的四个层级很多人一提“实时”默认就是毫秒级。其实数据新鲜度至少有四个层级每个层级的技术栈和成本差异巨大层级数据新鲜度代表方案成本T1离线24小时以上Hive/Spark离线调度低准实时10分钟~1小时短周期调度增量同步中低近实时1~5分钟微批流处理中实时秒级~毫秒级Flink/事件驱动/CEP高这里要特别说一句把“准实时”做好对大多数业务已经足够。很多报表和Dashboard做不到秒级的原因不是没有好引擎而是需求本身不需要。过度设计是数据团队最容易犯的毛病后文单独讲。2.3 实时的共同本质从“事后复盘”变成“事中决策”把上面几类需求放一起看会发现它们有一个共同点把数据从“事后复盘”变成“事中决策”。批处理解决的是“昨天发生了什么”实时解决的是“正在发生什么、要不要马上做点什么”。理解了这一点很多技术选型就不再纠结。比如一个场景如果业务根本不需要在事件发生的瞬间做决策那你完全没必要引入Flink用秒级或分钟级的定时任务就够了。真正的实时需求往往都直接指向决策动作——拦截一笔交易、触发一个告警、推送一条消息、调整一次库存。没有“决策”这个动作的实时多半是在自嗨。3. 从T1到T0架构演进的四个关键跳变3.1 第一跳从“全量重算”到“增量同步”早期想把数据实时化很多团队的第一反应是“把同步频率调高”。但全量同步在数据量上来之后必然扛不住几千万行每天全量扫一遍数据库压力大到业务方投诉。于是CDCChange 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进FlinkFlink按实体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或数仓增量目录再算一份日级结果两条链路共用同一份口径定义文件。日常查询走实时链路满足秒级场景日跑结果用于对账和修正发现差异后以离线口径为基准反推问题所在。这样既保证业务体验又给数据质量留了一条退路。对比维度LambdaKappa混合折中逻辑复杂度高两套代码低一套代码中共享口径口径一致性难保证天然一致靠对账兜底历史重放能力强依赖消息队列留存量中运维成本高中中高适合场景重口径、重稳定从零建设、小团队有存量资产的演进型团队我个人推荐有存量批处理资产的团队直接走第三条路经验是“演进”而不是“革命”。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%用准实时甚至T1就能解决。6.3 按数据量、延迟需求、口径复杂度做决策我常用的一个简易决策框架是这样的数据量小、延迟需求大于5分钟不需要引入流计算引擎用定时任务加增量同步即可。数据量大、但可以接受分钟级用微批或短周期调度别为“伪实时”增加过多成本。数据量大、要求秒级且强状态处理才考虑完整的Flink流处理链路。在线推理需求实时特征服务是刚需但特征计算可以先用规则做再逐步上复杂模型。这个框架不是精确公式但它能帮团队避开一个常见误区拿实时引擎去解决一个根本不需要实时的问题。6.4 我的个人经验与建议按我自己的实践体会批量到实时的演进最关键的是“需求驱动”而不是“技术驱动”。别因为看到别人都在搞Flink就觉得自己落后了先问清楚业务方到底需要多新鲜的数据需要哪个环节做决策能接受多少延迟和成本。如果确实要做也建议走增量演进先做增量同步和准实时链路跑通后再平滑升级到秒级流处理。每次演进都要保留离线批处理这条安全通道直到实时链路稳定运行几周。我自己踩过最深的坑就是过早地把离线任务停掉结果实时链路一故障全公司报表直接瘫痪最终花了整整一周才把离线口径重新对齐。从那以后我立了一条规矩实时链路要上线离线任务必须继续跑一段时间直到对账连续多日无差异才允许逐步退坡。这些经验不是教科书上教的是拿真实故障换来的。希望对做数据开发、数据架构的同行有所帮助。