资讯动态

Apache Spark Structured Streaming 迁移指南:从 2.4 到 4.4 的版本升级要点与配置兼容性解析

发布时间:2026/9/20 23:46:07 来源:尧图企业网站定制
大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载导读本指南完整梳理 Apache Spark Structured Streaming 从 2.4 到 4.4 各版本升级过程中的行为变更、弃用项与破坏性改动覆盖Trigger.Once弃用、Trigger.AvailableNow回退机制、checkpoint 元数据校验、AQE 对无状态流查询的支持、有状态算子分组键强制哈希分区等关键变更。文中每个升级要点均结合当前仓库中的 SQLConf.scala 配置定义与 StreamingErrors.scala 等源码实现进行纵深解析。读完本文你将能够评估升级影响面、正确配置回退开关、诊断 checkpoint 恢复失败并平稳完成流式作业的跨版本迁移。注意本迁移指南仅描述 Structured Streaming 专属的变更项。许多 SQL 迁移项同样适用于将 Structured Streaming 升级到更高版本请同时参考 Migration Guide: SQL, Datasets and DataFrame。一、从 Structured Streaming 4.3 升级到 4.4EXCEPT 与流式左输入的约束自 Spark 4.4 起使用EXCEPT且左输入为流式数据的流式查询会被直接拒绝。原因是优化器的重写rewrite会在最初的不支持操作检查之后引入流式聚合从而绕过原有校验可能导致查询语义与预期不符。影响范围涉及EXCEPT集合操作的流式查询。已运行的存量查询如果此前已通过检查点运行此类查询仍可从持久化 checkpoint 重启。这是因为 offset log 中记录了兼容性值compatibility valueSpark 会依据该值维持旧行为。实践建议升级前先审查所有流式查询的物理计划确认是否存在EXCEPT左输入场景升级后对存量查询应先在测试环境从原 checkpoint 验证重启行为再迁移到生产。二、从 Structured Streaming 4.1 升级到 4.2checkpoint 元数据文件缺失校验自 Spark 4.2 起从 checkpoint 重启流式查询时若满足以下条件将直接失败并抛出STREAMING_CHECKPOINT_MISSING_METADATA_FILEmetadata 文件缺失但offset 或 commit 日志中仍有数据。在旧版本中该场景会静默生成一个新的 query ID这在恰好一次exactly-once语义的 sink 中会造成数据重复——这是本变更被引入的核心原因。该错误条件在源码中有明确的单一定义点见 StreamingErrors.scaladef missingMetadataFile(checkpointLocation: String): Throwable { new SparkRuntimeException( errorClass STREAMING_CHECKPOINT_MISSING_METADATA_FILE, messageParameters Map(checkpointLocation - checkpointLocation) ) }恢复与兼容开关出现该错误时有两种处理路径恢复 metadata 文件推荐从备份或其他副本还原缺失的metadata文件使用全新的 checkpoint 位置重新消费并计算全部数据。若要恢复旧行为静默生成新 query ID可设置如下配置为false。该配置在 SQLConf.scala 中定义版本为 4.2.0默认truespark.sql.streaming.checkpoint.verifyMetadataExists.enabledfalse源码中的配置文档明确说明其动机当 offset 或 commit 日志包含数据时校验 checkpoint metadata 文件是否存在。这可以防止在 checkpoint 数据已存在时生成新的 query ID从而避免恰好一次 sink 中的数据重复。SQLConf 文档原文见上述文件链接。此变更对应 SPARK-55058。⚠️ 关闭该校验会重新引入数据重复风险仅在充分理解后果且确有兼容需求时使用。三、从 Structured Streaming 4.0 升级到 4.1AQE 支持无状态流查询自 Spark 4.1 起自适应查询执行AQE开始支持无状态stateless流式负载且默认开启因此升级后查询行为可能发生变化。总体而言 AQE 有助于提升性能例如**解决数据倾斜分区skewed partition**问题。该能力由内部配置spark.sql.adaptive.streaming.stateless.enabled控制定义于 SQLConf.scala# 4.1.0 引入默认 true属于 internal 配置 spark.sql.adaptive.streaming.stateless.enabledtrue值得注意的是源码文档强调启用该配置的前提是spark.sql.adaptive.enabled也必须开启To enable this config, spark.sql.adaptive.enabled needs to be also enabled.即两者为叠加关系。遇到回归如何回退如果升级后发现性能回退regression可显式关闭以恢复旧行为spark.sql.adaptive.streaming.stateless.enabledfalse实践建议升级到 4.1 后先在低流量环境运行无状态流作业如过滤、投影、无聚合的流式 ETL对比关闭 AQE 前后的吞吐与延迟指标再决定是否保留默认开启状态。四、从 Structured Streaming 3.5 升级到 4.0三项行为变更4.0 版本引入了三项与正确性、checkpoint 维护和路径解析相关的重要变更。4.1 Trigger.AvailableNow 的降级回退SPARK-45178自 Spark 4.0 起如果查询中任何 source 不支持Trigger.AvailableNowSpark 会回退到单批次执行single batch execution以避免 source 与 wrapper 实现之间不兼容可能导致的正确性问题、数据重复和数据丢失。该回退逻辑实现在 MicroBatchExecution.scala 的getTrigger()中当触发器为AvailableNowTrigger时Spark 会检查所有 source 是否实现SupportsTriggerAvailableNow接口——全部支持则使用MultiBatchExecutor()多批次否则打印警告并降级为SingleBatchExecutor()单批次case AvailableNowTrigger if (sparkSession.sessionState.conf.getConf( SQLConf.STREAMING_TRIGGER_AVAILABLE_NOW_WRAPPER_ENABLED)) { MultiBatchExecutor() } else { val supportsTriggerAvailableNow sources.distinct.forall { src val supports src.isInstanceOf[SupportsTriggerAvailableNow] if (!supports) { logWarning(logsource [${MDC(LogKeys.SPARK_DATA_STREAM, src)}] does not support logTrigger.AvailableNow. Falling back to single batch execution. ...) } supports } if (supportsTriggerAvailableNow) MultiBatchExecutor() else SingleBatchExecutor() }日志中同时提示降级为单批次执行时如果存在未提交的 batch可能无法保证处理到全部新数据建议与数据源开发者沟通使其支持Trigger.AvailableNow。4.2 checkpoint 目录额外空间比例配置SPARK-48931自 Spark 4.0 起新增配置spark.sql.streaming.ratioExtraSpaceAllowedInCheckpoint默认0.3用于控制 checkpoint 目录中允许的额外空间比例以存放维护任务maintenance task批量删除所需的过期版本文件从而摊销云存储上的 list 操作成本。该配置定义于 SQLConf.scala# 4.0.0 引入默认 0.3属于 internal 配置 spark.sql.streaming.ratioExtraSpaceAllowedInCheckpoint0.3其语义是批量删除时保留的最少过期版本数量由minBatchesToRetain * ratioExtraSpaceAllowedInCheckpoint计算得出。该计算式实现在 StateStoreConf.scalaMath.round(sqlConf.ratioExtraSpaceAllowedInCheckpoint * sqlConf.minBatchesToRetain)其中spark.sql.streaming.minBatchesToRetain默认值为100见 SQLConf.scala用于控制 checkpoint 文件state、offset、commit log的生命周期超过该批次数量的文件可被清理。因此默认配置下checkpoint 中至少会保留round(0.3 × 100) 30个版本的过期文件供批量删除。设置为0即恢复旧行为逐批删除。实践建议在对象存储S3、OSS、GCS 等上运行时保持默认值可显著减少 list 请求次数、降低成本在本地 HDFS 场景可按需调小。4.3 DataStreamWriter 相对路径解析提前到 DriverSPARK-50854自 Spark 4.0 起DataStreamWriter输出数据使用相对路径时绝对路径的解析在 Spark Driver 完成不再推迟到 Executor。这是为了让 Structured Streaming 的行为与 DataFrame APIDataFrameWriter保持一致。实践建议升级后检查所有使用相对路径作为 sink 输出的流式查询确认 Driver 端可解析到正确的绝对路径避免出现不同 Executor 解析不一致的问题。五、从 Structured Streaming 3.3 升级到 3.4Trigger.Once 弃用与 Kafka offset 获取默认值变更5.1 Trigger.Once 弃用SPARK-39805自 Spark 3.4 起Trigger.Once被弃用官方鼓励迁移到Trigger.AvailableNow。在当前的 4.x 源码中OneTimeTrigger仍然存在并被映射为单批次执行器见上文 MicroBatchExecution.scala 中case OneTimeTrigger SingleBatchExecutor()但新代码应使用Trigger.AvailableNow。迁移示例Scala// 旧写法已弃用 df.writeStream.trigger(Trigger.Once()).start() // 新写法 df.writeStream.trigger(Trigger.AvailableNow()).start()对应 Python 写法# 旧写法已弃用 df.writeStream.trigger(onceTrue).start() # 新写法 df.writeStream.trigger(availableNowTrue).start()5.2 Kafka offset 获取默认值反转影响 ACL 配置自 Spark 3.4 起配置spark.sql.streaming.kafka.useDeprecatedOffsetFetching的默认值从true变为false。该配置定义于 SQLConf.scala当前版本默认值确认为false# 3.1.0 引入3.4 起默认值由 true 改为 false spark.sql.streaming.kafka.useDeprecatedOffsetFetchingfalse默认不再依赖基于 consumer group 的调度方式获取 offset这会影响所需的 ACL 配置——使用新机制基于AdminClient时需要为查询主体配置对应的 Kafka 集群 ACL如 Describe、Read 等权限。详细说明请参考 Structured Streaming Kafka Integration 中的 Offset Fetching 章节。六、从 Structured Streaming 3.2 升级到 3.3有状态算子强制哈希分区SPARK-38204自 Spark 3.3 起所有有状态算子stateful operators都要求使用精确分组键exact grouping keys进行哈希分区。在旧版本中除 stream-stream join 外的有状态算子允许松散的分区标准loose partitioning criteria这为正确性问题留下了隐患。新版本通过强制精确分组键消除该风险。向后兼容机制对于由旧版本构建的 checkpointSpark 保留旧行为保证存量查询可以继续运行。这意味着存量查询从旧 checkpoint 重启行为不变新查询必须以精确分组键进行哈希分区否则可能失败或报错。实践建议升级后新建的聚合、去重、流式 join 等有状态查询需确保groupBy/join的键在分区语义上完全一致。七、从 Structured Streaming 3.0 升级到 3.1状态算子正确性检查与 Kafka offset 获取机制7.1 迟到行late rows导致的正确性问题检查在 Spark 3.0 及之前对于包含有状态操作、且可能向下游有状态操作发出早于当前 watermark 允许的迟到延迟的行即在下游有状态操作视角下的迟到行这些行可能被丢弃的查询Spark 只打印警告。自 Spark 3.1 起Spark 会主动检查这类存在正确性隐患的查询并默认抛出AnalysisException。该行为由配置spark.sql.streaming.statefulOperator.checkCorrectness.enabled控制定义于 SQLConf.scala3.1.0 引入# 默认 true检查并抛出 AnalysisException spark.sql.streaming.statefulOperator.checkCorrectness.enabledtrue若用户充分理解正确性风险、仍决定运行该查询可显式关闭检查spark.sql.streaming.statefulOperator.checkCorrectness.enabledfalse该问题根因与全局 watermarkglobal watermark机制有关多个有状态算子共享全局 watermark 时行可能在下游被过早丢弃。关闭检查前请务必完成正确性评估。7.2 Kafka offset 获取新机制AdminClient 替代 KafkaConsumer在 Spark 3.0 及之前Spark 使用KafkaConsumer获取 offset可能在 Driver 端造成无限等待infinite wait。自 Spark 3.1 起新增配置spark.sql.streaming.kafka.useDeprecatedOffsetFetching当时默认true可设为false以使用基于AdminClient的新 offset 获取机制。如上文所述该配置在 3.4 起默认值已改为false即当前版本默认使用AdminClient机制。相关说明同样参见 Structured Streaming Kafka Integration 的 Offset Fetching 章节。八、从 Structured Streaming 2.4 升级到 3.0schema 可空性、outer join 状态结构与 API 移除3.0 是迁移幅度较大的一次升级包含一个 schema 行为变更、一个 checkpoint 破坏性变更和三个 API 变更。8.1 文件类数据源 schema 强制可空forceNullable自 Spark 3.0 起通过spark.readStream(...)使用 text、json、csv、parquet、orc 等文件类数据源时Structured Streaming强制将 source schema 转为可空nullable。此前版本遵循 source schema 中的可空性标注但可能引发难以调试的 NPE空指针异常。该行为由配置spark.sql.streaming.fileSource.schema.forceNullable控制定义于 SQLConf.scala# 3.0.0 引入默认 true spark.sql.streaming.fileSource.schema.forceNullabletrue源码文档说明当为true时强制流式文件源的 schema包括所有字段为可空否则 schema 可能与实际数据不兼容导致数据损坏corruptions。要恢复旧行为可设置spark.sql.streaming.fileSource.schema.forceNullablefalse8.2 stream-stream outer join 的 state schema 变更SPARK-26154Spark 3.0 修复了 stream-stream outer join 的正确性问题改变了 state 的 schema。这带来一个关键的 checkpoint 兼容性影响如果从Spark 2.x 构建的 checkpoint启动使用 stream-stream outer join 的查询Spark 3.0 会直接失败若要重新计算结果必须丢弃 checkpoint 并重放之前的输入。实践建议升级到 3.0 时涉及 outer join 的流式作业应规划一次完整重算窗口而不是尝试复用旧 checkpoint。8.3 三个 API 变更变更类型旧 API已移除/隐藏新 API类移除org.apache.spark.sql.streaming.ProcessingTimeorg.apache.spark.sql.streaming.Trigger.ProcessingTime类移除org.apache.spark.sql.execution.streaming.continuous.ContinuousTriggerTrigger.Continuous类隐藏org.apache.spark.sql.execution.streaming.OneTimeTriggerTrigger.Once以ProcessingTime为例新旧写法对照// 旧写法类已移除编译失败 // import org.apache.spark.sql.streaming.ProcessingTime // df.writeStream.trigger(ProcessingTime(1000)).start() // 新写法 import org.apache.spark.sql.streaming.Trigger df.writeStream.trigger(Trigger.ProcessingTime(1000)).start()九、迁移要点速查表以下表格汇总本指南涉及的全部配置项便于升级前后统一核查配置项引入版本当前默认值说明spark.sql.streaming.checkpoint.verifyMetadataExists.enabled4.2.0truecheckpoint 有 offset/commit 数据时校验 metadata 文件是否存在内部配置spark.sql.adaptive.streaming.stateless.enabled4.1.0true无状态流查询启用 AQE需同时开启spark.sql.adaptive.enabled内部配置spark.sql.streaming.ratioExtraSpaceAllowedInCheckpoint4.0.00.3checkpoint 批量删除时保留的额外过期版本比例为0恢复旧行为内部配置spark.sql.streaming.minBatchesToRetain2.1.1100checkpoint 文件state/offset/commit保留的最小批次数spark.sql.streaming.kafka.useDeprecatedOffsetFetching3.1.0false3.4 起使用旧KafkaConsumer方式获取 Kafka offset内部配置spark.sql.streaming.statefulOperator.checkCorrectness.enabled3.1.0true检查全局 watermark 可能丢弃迟到行的正确性隐患spark.sql.streaming.fileSource.schema.forceNullable3.0.0true文件类数据源 schema 强制可空上述配置项多数标记为internal内部配置即不在官方文档默认列出但均可通过spark-defaults.conf或spark-submit --conf设置所有默认值与版本号均以当前仓库 SQLConf.scala 中的定义为准。升级路径建议跨大版本升级尤其 2.x → 3.x、3.x → 4.x时遵循逐版本核对 → 小流量试点 → checkpoint 兼容性验证 → 全量切换的流程涉及 outer join、EXCEPT与 state 相关变更时务必提前规划 checkpoint 重建与数据重算窗口。结合 Structured Streaming Programming Guide 与 SQL Migration Guide 一起评估可最大限度降低升级风险。赞分享大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载相关推荐Apache Spark SparkR 迁移指南从 1.5 到 4.0 的版本升级兼容性要点全解析Apache Spark SparkR 迁移指南从 1.5 到 4.0 的版本升级兼容性要点全解析 本指南以 Apache Spark 仓库中 docs/sp大数据数据分析批处理流处理机器学习图计算Apache Spark MLlib 迁移指南从 0.9 到 4.0 的完整升级路线与实践要点Apache Spark MLlib 迁移指南从 0.9 到 4.0 的完整升级路线与实践要点 本指南以 Apache Spark 当前仓库的 MLlib 迁大数据数据分析批处理流处理机器学习图计算Apache Spark Structured Streaming 与 Kafka 0.10 集成指南读写、偏移量与安全配置全解析Apache Spark Structured Streaming 与 Kafka 0.10 集成指南读写、偏移量与安全配置全解析 导读 本文是 Apach大数据数据分析批处理流处理机器学习图计算上一篇k-skill 复合技能实战biz-health-check 一查六源的韩国企业尽调数据交叉查询指南下一篇Reddit数据处理架构PostgreSQL与Cassandra的双引擎策略创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价