ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

Apache Spark Structured Streaming 状态数据源(State Data Source)集成指南:从 Checkpoint 读取状态与状态元数据

Apache Spark Structured Streaming 状态数据源(State Data Source)集成指南:从 Checkpoint 读取状态与状态元数据 大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载导读State Data Source 是 Apache Spark 4.0 引入的实验性数据源允许用户通过独立的批查询直接读取 Structured Streaming 查询 checkpoint 中保存的状态存储State Store键值对以及对应的状态元数据。它主要解决两类实际问题一是在测试中同时校验输出结果与状态内部数据二是排查有状态流式查询的故障通过查看状态追溯错误输出的成因。读完本文你将掌握statestore与state-metadata两种格式的完整用法、全部可选参数的语义与约束以及它们与transformWithState、stream-stream join 等有状态算子的配合方式。实验性声明本文介绍的 State Data Source 目前标记为实验性Experimental数据源选项与输出行为在后续版本中可能发生变化。截至本仓库所对应的 Spark 4.0 版本该数据源仅提供批查询读取能力写入write功能在未来的 roadmap 上。一、功能概述从 checkpoint 读取状态键值对State Data Source 的核心能力是通过运行一个独立的批查询从 checkpoint 中的状态存储读取键值对。在绝大多数情况下用户可以读取与单个有状态算子stateful operator匹配的单个状态存储实例即拿到该算子状态中的全部键值对。需要特别说明的例外是stream-stream join它内部会使用多个状态存储实例来分别缓冲左右两路输入。State Data Source 对此做了抽象封装用户无需关心内部细节即可按左/右侧读取状态详见下文读取 stream-stream join 的状态一节。从源码实现看该数据源由 StateDataSource.scala 实现它同时实现了TableProvider与DataSourceRegister注册的短名为statestore见shortName()方法。读取流程大致为解析StateSourceOptions→ 读取 checkpoint 的状态元数据与 schema → 构造StateTable→ 由分区读取器StatePartitionReader等逐分区扫描状态存储。二、以批查询读取状态存储全部默认参数2.1 三种语言的读取代码以下代码用默认参数读取 checkpoint 中某个状态存储的完整键值对Python / Scala / Java 三种 APIdf spark \ .read \ .format(statestore) \ .load(checkpointLocation)val df spark .read .format(statestore) .load(checkpointLocation)DatasetRow df spark .read() .format(statestore) .load(checkpointLocation);其中checkpointLocation是流式查询 checkpoint 的根目录路径。2.2 输出 Schema数据源返回的每一行包含以下三列列名类型说明keystruct取决于状态键的类型状态的键valuestruct取决于状态值的类型状态的值partition_idint状态存储所属的分区 ID需要注意key 和 value 的嵌套列结构高度依赖有状态算子的输入 schema 以及算子类型。因此强烈建议在读取前先用df.schema()或df.printSchema()查看实际的输出结构再据此编写后续的数据处理逻辑。从源码看该输出结构在 SchemaUtil.scala 的getSourceSchema()中生成默认非 change feed、非 transformWithState、非 join情况下输出key、value、partition_id三列其中partition_id为IntegerType。2.3 必选选项选项值含义pathstring指定 checkpoint 的根目录。既可以通过option(path, path)指定也可以直接用load(path)传入源码中path是唯一被强制校验的必选项在 StateDataSource.scala 的StateSourceOptions.apply()中若未提供path会直接抛出requiredOptionUnspecified(PATH)错误对应测试ERROR: path is not specified。2.4 可选配置项全表以下是statestore数据源支持的可选选项含义与默认值如下选项值默认值含义batchId数值最新已提交批次latest committed batch目标读取批次。用于时间旅行time-travel场景该批次必须已提交且尚未被清理operatorId数值0目标读取算子。当查询使用了多个有状态算子时使用storeNamestringDEFAULT目标状态存储名称。当有状态算子内部使用多个状态存储实例时使用除 stream-stream join 外一般不需要joinSidestringleft 或 right无目标读取侧。用于读取 stream-stream join 的状态snapshotStartBatchId数值无若指定则强制从该批次 ID 对应的快照开始读取随后重放 changelog 直至batchId或其默认值。注意快照批次 ID 从 0 开始等于快照版本 ID 减 1。必须与snapshotPartitionId一起使用snapshotPartitionId数值无若指定仅读取该特定分区。注意分区 ID 从 0 开始。必须与snapshotStartBatchId一起使用readChangeFeedbooleanfalse若为 true则读取状态在微批次间的变化change feed输出 schema 也会不同详见读取状态在微批次间的变化一节。该选项必须配合changeStartBatchId使用且不能与batchId、joinSide、snapshotStartBatchId、snapshotPartitionId同时使用changeStartBatchId数值无change feed 模式下第一个要读取的批次。需要readChangeFeed为 truechangeEndBatchId数值最新已提交 batchIdchange feed 模式下最后一个要读取的批次。需要readChangeFeed为 truestateVarNamestring无本次批查询要读取的状态变量名。如果使用了transformWithState算子则该选项必填。注意目前该选项仅适用于transformWithState算子readRegisteredTimersbooleanfalse若为 true可读取transformWithState算子中注册的定时器timer。目前仅适用于transformWithState算子且与stateVarName互斥同一时间只能使用其中一个flattenCollectionTypesbooleantrue若为 truelist state、map state 等状态变量的集合类型会被展平输出若为 false则值以 Spark SQL 的 Array 或 Map 类型返回。目前仅适用于transformWithState算子2.5 选项校验规则源码级补充这些选项并非简单的透传StateSourceOptions 中内置了完整的校验逻辑理解它们有助于避免运行时报错operatorId不能为负数invalidOptionValueIsNegative(OPERATOR_ID)。storeName不能为空字符串joinSide的值必须是left/right并且storeName与joinSide不能同时指定两者冲突时抛conflictOptions错误。snapshot 系列snapshotStartBatchId与snapshotPartitionId必须成对出现只给其中一个会抛requiredOptionUnspecifiedsnapshotStartBatchId不能为负且不能大于目标batchIdsnapshotPartitionId不能为负。change feed 系列readChangeFeedtrue时必须提供changeStartBatchIdchangeStartBatchId不能为负changeEndBatchId不能小于changeStartBatchIdreadChangeFeed不能与batchId、joinSide、snapshotStartBatchId、snapshotPartitionId同时使用反之在非 change feed 模式下指定changeStartBatchId/changeEndBatchId也会报错。transformWithState 系列readRegisteredTimerstrue与指定stateVarName互斥两者均为 boolean/string 类型非布尔值会抛错。读取方式限制该数据源仅支持批查询使用spark.readStream.format(statestore)会直接报错对应测试ERROR: trying to read state data as stream因为它只实现了批读取能力。以上绝大多数约束在 StateDataSourceReadSuite.scala 与 StateDataSourceChangeDataReadSuite.scala 中均有对应测试用例例如operator ID specified to negativebatch ID specified to negativestore name is emptyinvalid value for joinSide optionboth options joinSide and storeName are specifiedsnapshotStartBatchId specified without snapshotPartitionId or vice versaspecify changeStartBatchId in normal modejoinSide option is used together with readChangeFeed等可作为选项语义的行为规范参考。三、读取 stream-stream join 的状态Structured Streaming 内部通过多个状态存储实例来实现 stream-stream join这些实例在逻辑上组成了缓冲左右两路输入行的存储区。由于对用户而言按侧思考更直观State Data Source 提供了joinSide选项用于读取 join 特定一侧的缓冲输入joinSide取值为left或right。同时为了支持直接读取内部的状态存储实例也允许使用storeName选项——但**storeName与joinSide不能同时指定**源码中的conflictOptions(Seq(JOIN_SIDE, STORE_NAME))校验。从源码实现看当指定joinSide时StateDataSource.scala 会通过StreamStreamJoinStateHelper.readKeyValueSchema()读取对应侧的 key/value schema另外如果 stream-stream join 使用了虚拟列族state format version 3指定storeName的内部读取会被改写为按状态变量名stateVarName读取见modifySourceOptions()这一内部行为对用户透明。小贴士当查询中存在多个有状态算子例如 stream-stream join 之后再接去重 deduplication时需要借助operatorId定位目标算子而想要读取某个有状态算子内部的特定状态存储实例例如 join 的某一侧则需要stateStoreName——具体如何使用可参见下文State Metadata Source一节。四、读取 transformWithState 的状态transformWithState是一个允许用户在批次间维护任意状态的有状态算子。要读取它的状态需要在 State Data Source 的读取查询中额外指定一些选项。核心约束有两点必须指定stateVarName。transformWithState允许在同一查询中使用多个状态变量但它们可能具有不同的复合类型与编码格式因此需要一个批查询一次只读取一个状态变量。指定stateVarName即可选中感兴趣的状态变量。从源码看如果目标算子是transformWithState却未提供stateVarName会抛出requiredOptionUnspecified(stateVarName)如果提供了stateVarName但目标算子不是transformWithState例如普通聚合算子同样会报错。这些检查位于StateDataSource.runStateVarChecks()。定时器读取将readRegisteredTimers设为true可以返回跨所有分组键grouping key注册的全部定时器。源码中会校验算子属性如果transformWithState算子处于TimeModeNone则注册定时器不可用读取会报错同时该选项与stateVarName互斥。复合类型变量的两种读取格式Flattened展平默认复合类型被展平为单独的列输出。Non-flattened非展平复合类型以 Spark SQL 中单列的 Array 或 Map 类型返回。可根据自身的内存需求选择更适合的格式。这两种格式由flattenCollectionTypes选项控制默认true。展平与否会影响输出列的布局具体差异可从 SchemaUtil.scala 的generateSchemaForStateVar()看到状态变量类型flattentrue默认flattenfalseValueStatekey, value, partition_idkey, value, partition_idListStatekey,list_element, partition_id每个元素一行key,list_valueArray 类型, partition_idMapStatekey,user_map_key,user_map_value, partition_id每个键值对一行key,map_valueMap 类型, partition_idTimerStatekey,expiration_timestamp_ms, partition_id读取定时器时适用同左以 MapState 为例源码中展平模式通过unifyMapStateRowPair()将内部存储的复合键拆成分组键与用户键非展平模式则利用状态存储中相同分组键连续存放的特性按分组键聚合成一个Map{user_key - value}再输出。五、读取状态在微批次间的变化Change Feed如果想了解状态存储随微批次的演变过程而不是某一特定微批次下状态存储的全量快照应使用readChangeFeed选项。注意并非所有有状态操作都支持 change feed。某些操作如范围删除 range deletion无法表示为单条键值变化记录会直接报错请根据错误提示确认是哪个操作不受支持。例如以下代码读取从批次 2 到最新已提交批次的状态变化df spark \ .read \ .format(statestore) \ .option(readChangeFeed, true) \ .option(changeStartBatchId, 2) \ .load(checkpointLocation)val df spark .read .format(statestore) .option(readChangeFeed, true) .option(changeStartBatchId, 2) .load(checkpointLocation)DatasetRow df spark .read() .format(statestore) .option(readChangeFeed, true) .option(changeStartBatchId, 2) .load(checkpointLocation);change feed 模式的输出 schema 与普通模式不同列名类型说明batch_idlong变化所属的批次 IDchange_typestring两种取值update与delete。update表示插入一个不存在的键值对或更新已有键的值delete 记录的value字段为 nullkeystruct取决于状态键的类型状态的键valuestruct取决于状态值的类型状态的值delete 记录为 nullpartition_idint状态存储所属的分区 ID该输出结构同样由 SchemaUtil.scala 的getSourceSchema()生成readChangeFeed分支对于transformWithState的 ListState / MapStatechange feed 下的列布局与上文表格一致例如 MapState 会输出batch_id, change_type, key, user_map_key, user_map_value, partition_id。从源码看change feed 的选项约束包括changeStartBatchId必填且不能为负changeEndBatchId默认取最新已提交批次且不能小于changeStartBatchIdreadChangeFeed与batchId、joinSide、snapshot 系列选项互斥。这些约束在 StateDataSourceChangeDataReadSuite.scala 中有完整测试覆盖。六、State Metadata Source先了解 checkpoint再决定怎么读在使用 State Data Source 查询状态之前用户通常需要先了解 checkpoint 的相关信息尤其是关于状态算子的信息checkpoint 中有哪些算子与状态存储实例、可用的批次 ID 范围等。为此Structured Streaming 提供了名为State Metadata Source的数据源用于提供来自 checkpoint 的状态相关元数据。重要前提该元数据是在使用 Spark 4.0 运行流式查询时构建的。由更低 Spark 版本创建的已有 checkpoint不包含该元数据无法使用本元数据源查询。必须先使用 Spark 4.0 指向现有 checkpoint 重新运行一次流式查询以构建元数据之后才能查询。用户还可以通过batchId选项获取某个时间点的算子元数据。6.1 创建 State Metadata Store 批查询df spark \ .read \ .format(state-metadata) \ .load(checkpointLocation)val df spark .read .format(state-metadata) .load(checkpointLocation)DatasetRow df spark .read() .format(state-metadata) .load(checkpointLocation);该数据源在 StateMetadataSource.scala 中实现注册短名为state-metadatashortName()方法实现TableProviderDataSourceRegister底层表能力为BATCH_READ只支持批读取。6.2 必选与可选选项必选选项选项值含义pathstring指定 checkpoint 的根目录。既可以通过option(path, path)指定也可以直接用load(path)传入可选选项选项值默认值含义batchId数值最新已提交批次若无则取 0可选用于获取该批次处的算子元数据从源码看batchId的默认值逻辑是优先取最后已提交批次若不存在任何已提交批次则回退为 0OperatorStateMetadataUtils.getLastCommittedBatch(...).getOrElse(0L)。6.3 输出 SchemaState Metadata Source 的 schema 是静态的源码中StateMetadataTableEntry.schema直接定义每行包含以下列列名类型说明operatorIdint算子 IDoperatorNamestring算子名称stateStoreNameint状态存储名称numPartitionsint分区数minBatchIdint可查询状态的最小批次 ID。注意若正在运行占用该 checkpoint 的流式查询该值可能失效因为清理cleanup会持续进行maxBatchIdint可查询状态的最大批次 ID。注意若正在运行占用该 checkpoint 的流式查询该值可能失效因为查询会继续提交更多批次operatorPropertiesstring算子使用的属性列表编码为 JSON。此处输出取决于具体算子_numColsPrefixKeyint元数据列隐藏列除非用 SELECT 显式指定才会出现其中_numColsPrefixKey的实现为MetadataColumn源码注释state store instance 前缀键prefix key的列数属于隐藏元数据列。6.4 典型用途该数据源的一个主要用途是当查询包含多个有状态算子时借助元数据识别operatorId。例如stream-stream join 后跟去重deduplication这种组合场景列operatorName可以帮助用户确定给定算子对应的operatorId进而用statestore源配合operatorId选项精确读取目标算子的状态。此外如果想读取某个有状态算子如 stream-stream join内部的特定状态存储实例列stateStoreName可用来确定目标——配合statestore源的storeName选项即可精确读取。七、实际使用建议与总结先元数据、后状态面对陌生 checkpoint 时先运行state-metadata查询确认可用算子operatorId/operatorName、状态存储实例stateStoreName、批次范围minBatchId/maxBatchId再决定statestore查询需要的batchId、operatorId、storeName/joinSide等参数能显著减少试错成本。先 Schema、后解析statestore的key/value嵌套结构高度依赖算子输入与算子类型务必先用printSchema()确认结构再编写字段访问逻辑对于transformWithState的集合类状态变量还需要注意flattenCollectionTypes带来的列布局差异。明确时间语义batchId实现时间旅行但目标批次必须已提交且未被清理change feed 模式readChangeFeedchangeStartBatchId则用于观察状态的逐步演变两者不可混用。版本与实验性前提State Data Source 与 State Metadata Source 均为 Spark 4.0 的实验性能力元数据只有在 Spark 4.0 运行过查询后才存在选项与输出行为可能随版本调整生产使用前请以当前版本仓库中的 structured-streaming-state-data-source.md 文档与对应源码、测试为准。从整体看这两个数据源把原本只能黑盒观察的流式状态存储变成了可查询的批数据表为有状态流式查询的可测试性与可观测性提供了直接支撑测试侧可以同时断言输出与内部状态运维侧可以在故障发生时回溯状态来源。随着未来写入write功能的加入State Data Source 有望在状态修复、状态迁移等场景发挥更大作用。赞分享大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载相关推荐DeerFlow 工具集成实战4 个配置让 AI 深度研究真正跑起来DeerFlow 工具集成实战4 个配置让 AI 深度研究真正跑起来 让 AI 做行业调研很多人第一步就卡在调不动搜索引擎上模型只会生成文字却拿不到人工智能大模型AI Agent自主智能体工具调用MCP 服务Agent 沙箱AI 技能后端前端Alpine.js State 状态管理完全指南x-data 本地状态与 Alpine.store 全局状态Alpine.js State 状态管理完全指南x data 本地状态与 Alpine.store 全局状态 Alpine.js当前仓库为 alpine 前端Rocket 状态管理实战指南Managed State、请求本地状态与数据库连接池Rocket 状态管理实战指南Managed State、请求本地状态与数据库连接池 导读 Rust Web 框架 Rocket 为应用提供了完整、线程安全的后端开发工具上一篇企业级ONNX模型部署架构设计5大核心策略与性能优化实战下一篇告别渐变兼容烦恼Autoprefixer一键修复linear-gradient语法创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表