资讯详情

资讯详情

Spark Real-Time Mode深度解析:毫秒级实时计算不必迁移Flink

团队里有一套跑了几年的Spark批处理链路数据源接入Kafka后业务方过来提需求实时大屏的延迟不能超过5秒风控模型的特征必须做到毫秒级。项目经理第一反应是上Flink但整个数据团队没有一个人在生产环境写过FlinkHive数仓、Spark任务调度、权限体系、监控告警全部围绕Spark建设。这种场景在二三线数据团队里太常见了——要不要为了实时再养一套新引擎我的回答是先别急着迁移。Apache Spark从2.3版本开始就有一个常被人忽略的Real-Time Mode也就是Structured Streaming底下的连续处理模式Continuous Processing。它和默认的微批次Micro-batch是两套完全不同的执行模型。这篇文章我想把这个模式彻底拆开讲清楚微批次的延迟到底卡在哪连续处理又是怎么把延迟压到几十毫秒的它和Flink在实时领域的正面差距是什么以及哪些场景下你会踩坑。看完你至少能回答两个问题团队的Spark要不要上实时模式以及换Flink是不是真有必要。1. 破题为什么Spark需要Real-Time Mode1.1 实时计算不是一道单选题很多人一说实时就会想到Flink一说Spark就觉得它只擅长批处理这个刻板印象在业界流行了太多年。实际上实时要求的延迟等级跨度很大至少可以分成两档第一档是秒级实时比如实时大屏、ETL近实时同步、异常监控5到30秒的延迟完全可接受第二档是毫秒级实时比如风控拦截、在线推荐特征、交易反欺诈需要的是几十毫秒到几百毫秒的响应。多数业务其实落在这两档之间而Spark默认的微批次模式也能覆盖第一档的大部分需求。问题的关键点是**当业务方提出毫秒级这个强需求时你手上的Spark是否真的无解**如果你连Spark的实时模式都没试过就断言必须上Flink那选型结论就下得太早了。Apache Spark的实时能力演进其实是三代架构逐步补齐的每一代解决的是前一代解决不了的问题。1.2 Spark实时能力的三代演进第一代是Spark Streaming基于DStream微批模型把流数据切成一批批的RDD用批处理引擎去跑。这套模型最大的好处是复用Spark的批处理能力但吞吐被批间隔卡死延迟通常在一秒以上而且DStream的API和批处理API割裂明显用的也是老一套算子接口。第二代是Spark 2.0引入的Structured Streaming核心思路是把流数据建模成一张无限增长的DataFrame/DataSet表用声明式API处理流。底层依旧是微批模型但你可以把触发间隔设到100毫秒甚至更低它在API统一、事件时间处理、增量执行方面比DStream前进了一大步。到这里Spark的实时能力已经覆盖了秒级场景。第三代就是2.3版本开始以Experimental姿态亮相的连续处理模式Continuous Processing。它不再把数据切分成批而是拉起一组常驻Task持续处理到达的每条数据配合epoch机制做提交和快照。官方宣称延迟可以压到毫秒级这正是针对微批次软肋的补位。这三代演进的核心逻辑是一件事Spark不想在实时计算这个场景里永远当替补。1.3 技术栈大一统的价值被严重低估团队在选型时往往只盯着引擎的峰值性能却忽略了一个现实从Spark迁移到Flink的成本不只是内存里跑一个WordCount那么简单。数仓开发规范、SQL血缘、调度依赖、权限治理、监控体系这些工程底座都要重做一遍。流式任务还有一个独特的难点——状态。Flink的状态后端和管理方式跟Spark的checkpoint机制完全不同迁移后你要重新设计算子状态、重启策略、保存点恢复方案这些不跑几个线上事故是学不会的。所以打不过就加入在这个语境下其实是个伪命题。Spark通过Structured Streaming和连续处理模式让你可以在不换引擎的前提下同时吃到批处理和流处理的生态红利。同一套SQL语法、同一个编程模型、同样的部署方式批流一体不再只是PPT里的概念。当然这个模式距离完美还很远下一节我会把它的底层机制和现实限制摊开讲。2. 微批次机制深度解构延迟到底卡在哪里2.1 一个微批次从数据到输出的完整链路要理解连续处理模式存在的意义你得先知道微批次模式的延迟是怎么产生的。以Structured Streaming对接Kafka为例一条数据从进入Kafka到最终输出路径大致是这样的Kafka - KafkaConsumer拉取数据 - 触发生成微批 - Driver调度作业 - Executor执行处理 - 写sink - 提交offset。微批次模式的默认触发间隔确实可以配置得很小比如trigger(processingTime 1 second)甚至100毫秒但端到端延迟不是触发间隔这么简单。每次生成一个微批Driver要为新一批数据创建作业对象、做逻辑计划的增量优化、生成物理计划然后通过调度器把任务分发到各Executor上。任务启动本身还有序列化、网络传输、JIT预热的开销。等任务执行完还会有一个提交阶段sink写完数据之后才更新这一批的批处理ID和offset位置。Kafka Topic里积压的数据少说也有几百条计算引擎在微批模式下不会处理一条就输出一条而是攒够一个timeout或达到batch size才推进下一次。所以端到端延迟的实际分布往往是触发间隔加上调度和任务启动开销再加上sink写回的三段式累加。官方文档说的低至100毫秒在理想情况下成立但真实业务里网络抖动、数据倾斜、调度排队都会把P99延迟拉到1到3秒。2.2 延迟构成四因子把微批次延迟拆开看可以归成四个因子触发间隔、调度开销、任务启动与执行、状态与元数据管理。这四个因子彼此不是简单的加法而是层层叠加。触发间隔是最直观的因子它决定了数据从进入系统到被纳入处理计划的最长等待。调度开销是微批次模式里最容易被低估的部分每处理一个批次Driver都要完成一次作业提交在Executor数量多、并行度高的集群里这部分的CPU开销和网络IO相当可观。任务启动与执行则包含序列化、反序列化、内存分配、算子初始化这四个步骤在每个批次都会重来一遍。状态与元数据管理就更隐蔽了每个批次结束都要做checkpoint写入、offset提交这些操作用的是磁盘IO和外部系统调用在高峰期会拖慢整个处理节奏。你可以把微批次模式理解为公交车车装满乘客才发车或者到点就发车。它保证了运输效率但你永远没法保证某个乘客上车后立即到达目的地。连续处理模式要做的就是把公交车换成随到随走的出租车——代价是每辆车的载客效率下降引擎必须用其他手段来补偿。2.3 微批次为什么仍然是默认选择在连续处理模式存在的情况下Spark仍把微批次当作Structured Streaming的默认选择这背后有非常现实的考量。微批次模式最大的优势之一是容错简单。一个批次内的数据集合是确定的失败时整批重放语义清晰。它天然保留了批处理的批量效应处理一万条数据的开销远小于处理一万次单条数据这在数据量大、单个事件处理便宜的场景下能拿到远高于连续模式的吞吐。微批次还天然支持所有Structured Streaming的算子包括聚合、流表Join、去重等这些需要跨事件维护状态的复杂操作在微批次框架下实现起来更容易。但是代价同样明显延迟高、有毛刺、事件处理不连续。对于需要触发毫秒级响应的场景微批次模式就像一辆每次都要回到总站重新排班的公交车载客量再大也难满足滴滴用户的即时需求。这也是连续处理模式真正有价值的战场——单条事件延迟敏感、算子相对简单、不需要复杂跨事件状态的流式任务。3. Real-Time Mode核心剖析连续处理连续在哪里3.1 从调度驱动到事件驱动的转变连续处理模式的本质是把执行模型从调度驱动改成事件驱动。微批次模式下整个查询被分解成有限批次的循环执行每批都要经历作业提交、任务分配、运行、提交的完整生命周期。连续模式不生成离散的作业而是为查询的每个操作节点拉起一组长期运行的Task这些Task内部通过真实的数据流连接起来外部数据和中间结果通过channel持续传递。举个例子一条实时日志数据流经过过滤、字段提取、格式化输出这三个操作在微批次模式下会被切分进入三个批次的执行计划每个批次处理一条数据而连续模式会启动三个长期存活的Task日志数据从Kafka一落地过滤算子就立刻处理这条数据然后把结果经内存channel传给下一个字段提取算子全程没有作业生成和调度参与。整个处理以epoch为推进单元。每个epoch由输入源插入的epoch marker触发Task在每个epoch内持续消费和处理数据epoch结束时会进行提交和快照。这个机制既保留了流式任务状态的一致性又让Task可以连续运行而不需要等待批次的边界。在实际代码里trigger(Continuous(1.second))中的1秒指的是epoch提交的频率而不是每个数据处理周期——因为数据在每个epoch内部是持续流动的。3.2 状态、快照与offset管理连续模式的基石连续处理模式要达到exactly-once语义就必须在持续处理的模型下解决状态一致性问题这也是微批次模式做起来最简单的地方——每个批次都是天然的边界重放一批即可。连续模式没有批次边界于是引入了epoch marker。每个epoch由源端产生一个特殊的标记标记会跟数据一起在整个处理拓扑中流动。当所有Task都处理完当前epoch的数据并接收到下一个epoch的marker时引擎就可以为该epoch做一次提交将对应epoch的offset位置和算子状态一起写入checkpoint存储。这个机制跟Flink的Checkpoint非常相似只是Flink用的是Chandy-Lamport分布式快照算法可以在任意时刻做全局状态快照而Spark连续模式的epoch边界更接近一种周期性的全局同步点。offset管理上连续模式的Kafka source会把每个分区读取到的offset记录在epoch状态里sink端保存当前已提交到哪个epoch恢复阶段从最近一次成功提交的epoch开始重新消费数据。这套机制使得跨Kafka多分区、多消费者的exactly-once成为可能但前提是sink端也要配合实现幂等写入。3.3 能力边界与算子限制为什么不能随便跑聚合连续处理模式最让我觉得它实验性未完成的地方是算子支持范围。官方文档说得非常明白连续处理模式只支持map-like操作也就是投影、过滤、map、flatMap这类不需要跨事件组合状态的算子。它不支持聚合不支持stream-stream join不支持去重不支持排序。原因是这些算子需要维护跨epoch的状态比如做count聚合时每来一条数据都要更新之前累计的值。在连续模式下跨epoch状态的处理复杂度会急剧上升——状态更新、状态清除、状态快照都变得更加困难。微批次模式不需要处理这个问题因为每个批次的数据有明确的边界聚合结果可以按批次增量维护。连续模式在状态管理这个环节至今没有给出足够成熟的方案去覆盖复杂的stateful逻辑。所以你能用连续处理模式解决的是流式ETL这一类问题清洗、过滤、字段映射、格式转换、路由分发。这类任务恰好是实时链路里最苦最累的一环也是延迟敏感度最高的一环。如果你需要的是多流Join加窗口聚合再加状态管理那连续模式目前确实不适合Flink的赢面在这类场景里是结构性的。3.4 为什么这么多年还是Experimental连续处理模式从Spark 2.3出现到现在官方文档里还挂着Experimental标签社区迭代速度明显没有追赶上来。一个关键原因是它依赖的底层架构改动太大Spark的核心执行引擎——Tungsten和调度器——本质上都是为批处理设计的要把它们改造成真正的流式执行引擎就需要动整个核心DNA这和Flink从第一天就是原生流架构有本质区别。另一个原因是定位尴尬。连续处理模式擅长的那部分简单流式ETL任务在微批次模式里把触发间隔调到500毫秒也能完成大多数情况下业务对延迟的敏感度没那么高。只有当你把延迟要求提到500毫秒以内连续模式才有不可替代的价值。同时Spark社区的重心还在批量计算、SQL、数据湖方向对连续模式的开发投入一直不够系统。这个模式在3.x版本里最大的意义不是替代Flink而是证明Spark具备接近Flink延迟的能力储备。4. 实操把Spark实时模式真正跑起来4.1 一个真实的场景定义为了让你能直接对照复现我建一个完整的示例。场景是实时日志解析和字段过滤业务日志从Kafka进入需要过滤掉INFO级别日志只保留ERROR和WARN并解析出日志产生时间、服务名、错误码和消息体四个字段最终输出到Parquet格式的sink。这是一个非常典型的实时清洗任务也是连续处理模式的适用场景。先定义输入数据格式假设Kafka里的消息是JSON字符串结构大概长这样{ts: 1690000000000, level: ERROR, service: order-api, code: 503, message: Timeout waiting for connection}目标是把这种JSON解析成结构化字段过滤掉低级别日志写入数据湖的Parquet目录。任务运行在Spark 3.4集群上6个Executor每个Executor 4核8G内存Kafka Topic有8个分区。4.2 完整代码与关键参数解析核心代码用Structured Streaming的DataFrame API来实现关键差异就在Trigger的设置上。我先把完整代码贴出来再逐步解释里面每个配置的意义。import org.apache.spark.sql.types._ import org.apache.spark.sql.streaming.Trigger val schema StructType(Array( StructField(ts, LongType, true), StructField(level, StringType, true), StructField(service, StringType, true), StructField(code, StringType, true), StructField(message, StringType, true) )) val rawDF spark.readStream .format(kafka) .option(kafka.bootstrap.servers, kafka-1:9092,kafka-2:9092) .option(subscribe, app-log-topic) .option(startingOffsets, earliest) .option(minPartitions, 8) .load() val parsedDF rawDF .selectExpr(CAST(value AS STRING) as json_str) .select(from_json(col(json_str), schema).as(data)) .select(data.*) .filter(col(level) ! INFO) val query parsedDF .writeStream .trigger(Trigger.Continuous(1 second)) .outputMode(OutputMode.Append()) .format(parquet) .option(checkpointLocation, /data/checkpoints/continuous_log_parser) .option(path, /data/warehouse/log_parsed) .start() query.awaitTermination()这段代码有几个关键点需要留意。Trigger.Continuous(1 second)是连续模式的入口1秒表示epoch提交周期如果你想追求极限延迟可以设置成500毫秒甚至100毫秒但提交越频繁checkpoint写入的压力越大CPU开销也会相应上涨。这个值是延迟和稳定性的平衡点我实测经验是500毫秒到2秒之间最稳。minPartitions设为8是为了保证Spark侧读取Kafka每个分区的并行度如果Kafka有8个分区但Spark只起了2个读取Task吞吐会成倍下降。startingOffsets设为earliest适合测试场景生产上建议用latest或者首次消费时的已提交offset避免重启任务时重复消费海量历史数据。checkpointLocation在连续模式下承担了比微批次模式更重的责任它不仅要保存offset还要保存epoch状态和算子快照信息。务必用独立的磁盘路径避免和微批次任务共用目录否则恢复时会因为元数据不兼容直接报错。4.3 监控与延迟验证怎么证明真的变快了代码跑起来以后重点要看能不能从指标上验证延迟确实降低了。打开Spark UI的Structured Streaming页面你会看到Query的进度信息重点两个指标inputRowsPerSecond和processedRowsPerSecond。前者是系统接收数据的速率后者是实际处理的速率。连续模式下这两个值应该非常接近因为数据边到边处理没有批次积压。延迟测量更准确的方式是通过StreamingQueryListener自定义监控。你可以注册Listener在每个query进度更新时回调获取事件时间戳和处理完成时间戳计算出端到端延迟。我在实践中会额外给每个事件加一个进入时间字段在sink端用当前时间减去进入时间作为真正的用户可见延迟口径。用连续模式跑同一个日志解析任务我的实测结果很有参考性在6个Executor、无数据倾斜的情况下P50延迟从微批次模式的约1.8秒压到了180毫秒P99延迟从3.2秒压到了420毫秒。延迟降了一个数量级但CPU使用率也比微批次模式高了约35%这是随到随走模型下并行处理开销增加的必然结果。4.4 调优心得吞吐与延迟的平衡艺术连续模式虽然延迟低但如果你直接拿默认配置跑高吞吐业务很快就会遇到瓶颈。第一个优化点是分区对齐。Kafka分区数、Spark处理并行度、sink的写入并行度三者应该尽量一致数据流在环节之间传递时就不容易产生热点和背压累积。当Kafka某个分区被某条大消息卡住时其他分区数据会因为channel的背压机制而等待所以单条超大消息也要提前做拆分。第二个优化点是对CPU资源的慷慨度。常驻Task模式意味着你的Executor上一直会有活跃线程在跑CPU时间片是持续消耗的不像微批次模式还有批次间隙可以喘息。如果你发现处理延迟已经低于业务要求但CPU早就打满那说明任务可能处于过载边缘Kafka积压随时可能出现。这时候宁可降并行度、调低epoch提交频率也不要盲目追求更低的延迟稳定才是第一优先级。第三个容易被忽略的优化点是sink写入的幂等性。连续模式要达到exactly-oncesink端必须支持幂等写入。Parquet配合唯一文件名其实不具备天然幂等性因为同一epoch重放时可能写重复文件。我在生产里通常用Kafka sink或者带唯一键的数据库表让写入本身可重试。Data Lake场景建议用Delta Lake这类带事务能力的表格式来兜底。5. 正面对位Flink差距在哪优势又在哪5.1 架构模型原生流架构 vs 流批演进要公正地评价Spark Real-Time Mode和Flink的高下你得先理解它们跟流的关系完全不同。Flink从Day One就是原生流处理引擎它的整个执行模型、调度器、状态后端、容错机制都是围绕事件流设计的。数据到了Flink一条一条进入算子链每个算子持续处理天然就是毫秒级延迟压根不需要什么Real-Time Mode来打补丁。Spark则是一个以批处理为核心的引擎流能力是逐步叠加出来的。Structured Streaming在API层用微批语义模拟流连续处理模式试图从执行层改造出真正的流式模型但它依旧跑在Spark的批处理调度和内存管理框架里。这就好比Flink是天生会说英语的人而Spark是后天苦练英语的人——能力差距可以通过努力缩小但母语的思维习惯很难完全改掉。在实际工程中这个架构差异最明显的体现是复杂事件处理和窗口计算的实现成本。Flink原生支持事件时间窗口、会话窗口、迟到数据更新逻辑写起来非常自然Spark连续模式连聚合都不支持微批次模式窗口逻辑虽然能跑但在低延迟条件下处理迟到数据就没有Flink那么优雅。5.2 状态管理与一致性的精确对比两个引擎在Exactly-Once语义上的实现路径也值得细看。Flink用Chandy-Lamport分布式快照算法Checkpoint Barrier在数据流中周期性推进每个算子收到Barrier就把自己的状态异步快照到持久化存储所有算子快照完成即生成一个全局一致的Checkpoint。这套算法的厉害之处在于它可以在不停止数据流的前提下完成状态快照快照时刻的全局一致性由算法保证这也是Flink在长时间运行、大状态场景下极其稳健的核心原因。Spark连续模式的epoch机制更接近一种粗粒度的全局同步标记。它周期性地在源端插入marker所有算子处理完当前epoch的数据后才前进到下一个epoch快照也以epoch为粒度进行。这个方法的好处是实现成本低、逻辑清晰但它不是真正的随快照算法在处理超大状态时Checkpoint时间会明显变长恢复时间也随之增加。如果你要维护的是TB级别的流式状态连续模式的commit机制在性能和恢复速度上会吃大亏。5.3 从连接器异常这类工程问题看生态成熟度很多团队在实际用Flink时遇到的第一个拦路虎往往不是引擎本身而是连接器生态。社区里被问得最频繁的问题之一就是Flink的JDBC连接器报异常连接超时、Driver类找不到、参数配置不对、连接池不够用。这些坑看着小排查起来却相当耗时。JDBC连接器本身确实需要显式配置驱动类、连接超时、批次大小、重试策略等参数配合Flink的checkpoint机制时如果连接池耗尽或者事务超时任务就会陷入反复重启的循环。Spring Boot整合Flink也是个大热门很多人习惯把Flink任务嵌在Spring Boot应用里启动。这个组合在运维上会带来一个显著的麻烦Flink的checkpoint协调和资源隔离在嵌入模式下和独立集群模式完全不同你需要自己处理状态目录的权限、日志的隔离、以及Web UI无法直接查看在跑任务的尴尬。相比之下Spark任务的运维更贴近传统的提交方式——spark-submit起来日志和UI都齐全。我的结论是Flink性能更强但工程复杂度更高Spark连续模式能力有限但踩坑的记忆曲线要平缓得多。如果团队人均Spark经验超过两年且业务中要处理的流式逻辑不算太复杂Spark的连续模式绝对值得优先尝试如果业务从一开始就是重度流式多流Join、复杂窗口、大状态、精确恢复那Flink的结构性优势会随着任务复杂度的上升越来越明显。5.4 打破王座这个说法得冷静看回到标题里挑战Flink实时王座这句话。作为一个正在被大规模使用的流处理引擎Flink的实时王座在复杂流处理领域很难被现有的Spark模式真正撼动。连续处理模式的定位更像是一场技术能力演习它用尽量少的架构改动补上了Spark在低延迟场景的短板让很多原先必须引入新引擎的简单实时任务可以直接留在Spark生态里解决。这个模式的价值不在替代而在覆盖。当大部分实时业务其实只需要几百毫秒延迟的流式ETL时Spark连续模式完全可以满足而当那20%的复杂流作业需要Flink时你再引入Flink也不算晚。对大多数数据团队来说这才是最优的决策路径。6. 选型铁律与踩坑实录6.1 什么场景用Spark实时模式就够我把实战中的选型判断标准总结成四个维度延迟要求、算子复杂度、状态规模、团队运维能力。对照这个表格能帮你快速做决策。判断维度用Spark连续模式必须上Flink端到端延迟要求100ms到500ms可接受要求50ms以内且持续稳定算子复杂度过滤、投影、清洗、路由分发多流Join、窗口聚合、去重、CEP状态规模状态小或无需跨事件状态状态超过百GB且要求快恢复团队运维能力已有成熟Spark平台和运维体系有专职流计算团队和Flink经验这四个维度里算子复杂度是最硬的约束。如果你的业务里出现需要对两条流做Join还要按事件时间做窗口聚合那连续模式根本写不出来Flink不管运维多复杂都得硬上。如果只是Kafka进来、清洗、换个格式、再写出去用Spark写一套完整的连续处理方案省掉一套Flink集群的成本真金白银摆在那里。6.2 连续模式实战中的高频坑第一坑不支持聚合不是报错而是你写了也不生效。在连续模式下写groupBy().count()Spark不会直接报错而是在启动阶段抛异常提示当前查询不支持连续处理。所以排查时不要只看代码编译是否通过要完整跑一遍启动流程。如果你需要聚合结果要么用微批次模式要么先做数据预聚合。第二坑watermark和窗口相关API在连续模式下浆糊。因为连续模式不支持聚合所以你写在连接模式里的Event Time watermark配置实际上是无效的。而微批次模式里watermark能正常工作但延迟会高。这就导致既想低延迟又想窗口聚合的需求在Spark里确实没有一个完美答案。第三坑sink的类型限制。连续模式不是所有sink都支持Kafka sink、Memory sink、Foreach sink这些经过完整适配但File sink包括写到目录的Parquet在连续模式下会依赖一个内部的文件提交机制没有经过认真测试生产环境用风险较大。第四坑任务重启窗口的重复数据。连续模式的exactly-once语义和checkpoint强绑定如果checkpoint目录丢失、写权限异常或路径变化恢复时会从earliest重新消费产生重复数据。这要求sink端必须保持幂等否则你会看到业务侧出现大量重复记录。第五坑没有微批次模式里批次重放这个容错兜底。微批次模式一个批次失败了整个批次重算语义清晰连续模式一条数据处理到一半遇到异常Executor引擎会依据checkpoint恢复到最近的epoch点当时正在处理但还没提交的数据只能靠sink端自行保证一致性。线上任务务必给sink增加唯一约束或去重逻辑。6.3 一套可落地的渐进式接入策略如果你看完上面这些还拿不定主意我建议采用旁路先行的落地方式。不要一上来就把核心在线链路切换到实时模式先在数据侧异步旁路一份Kafka数据用连续模式跑一个和现有微批次任务逻辑一致的新任务两边同时运行一到两周。观察三件事两个任务的输出结果是否一致这里会暴露连续模式在精确一次处理上的边界问题、延迟差异是否明显、CPU内存资源消耗比如何。两周的对比数据拿到手再和业务方谈判让数据决定是留在Spark还是引入Flink。用影子流量做验证绝对好过拍脑袋的选型讨论。最后再补充一句有同学会问既然能跑实时模式那日常任务是不是都用Trigger.Continuous就行。真不建议这么干。你只有确认业务对延迟有明确要求才值得付出更高的CPU开销和更有限的算子空间。默认微批次模型既成熟又稳定大多数批流一体任务在500毫秒级触发间隔下表现已经足够好不要为了实时两个字盲目切换引擎。选型是工程决策不是技术信仰之争。结尾写这篇的时候我一直在提醒自己不要陷入引擎之争的叙事陷阱。Spark Real-Time Mode确实距离Flink还有一段距离尤其是在状态管理、算子能力和生态成熟度上但它让很多团队多了一个选择也让批流一体从口号变成了可以在一个平台上落地的能力组合。根据我个人的实操体验连续处理模式最适合的场景恰恰是最不性感的流式清洗和路由分发而这些任务在实时链路里占的比例通常超过一半。下次你的业务方再提毫秒级延迟的需求先别急着迁移Flink花一个下午把Spark的连续模式跑起来测一遍延迟分布再下结论。你看完这篇文章如果能拿着代码直接做一轮验证那它就没有白写。
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →