大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载导读本文以 Apache Spark 官方 PySpark API 参考中的Structured Streaming Input/Output输入输出主题为骨架系统讲解流式数据读写两大核心接口DataStreamReader流式读取与DataStreamWriter流式写出的全部公开方法。你将掌握从文件系统、Kafka 等外部存储读取数据流、通过多种触发器与输出模式把结果写入各类 Sink、以及基于foreach/foreachBatch实现任意下游处理的完整实战能力并了解每个 API 在 readwriter.py 中的底层实现与参数语义。本文对应官方文档入口为 python/docs/source/reference/pyspark.ss/io.rst其属于 python/docs/source/reference/pyspark.ss/index.rst 中 Structured Streaming API 总览的 Input/Output 一章DataStreamReader与DataStreamWriter的完整实现位于 readwriter.py1964 行类中每个方法都带有可直接运行的 doctest 示例可在SPARK_HOME下通过python -m pyspark.sql.streaming.readwriter自测。一、概览流式 I/O 的两个入口Structured Streaming 的 I/O 由两个对称的接口构成均自 2.0.0 引入并从 3.5.0 起支持 Spark ConnectDataStreamReader负责把外部存储系统文件系统、键值存储等中的数据加载为流式DataFrame通过SparkSession.readStream获得DataStreamWriter负责把流式DataFrame的输出写到外部存储系统或自定义 Sink通过DataFrame.writeStream获得。# 读取端入口 spark.readStream # - DataStreamReader # 写出端入口 df.writeStream # - DataStreamWriter入口方法的源码位置session.py 的 readStream、dataframe.py 的 writeStream。这两个接口正在演进evolving官方标注为实验性 API使用时需注意后续版本可能调整。io.rst 通过 autosummary 一共列出 20 个公开方法其中DataStreamReader13 个、DataStreamWriter11 个format/option/options两接口共有接口方法DataStreamReaderchangescsvformatjsonloadnameoptionoptionsorcparquetschematabletextDataStreamWriterforeachforeachBatchformatoptionoptionsoutputModepartitionByqueryNamestarttoTabletrigger下面按读取 → 写出 → 触发与提交 → 进阶 Sink的链路逐一展开。二、DataStreamReader把外部数据读成流式 DataFrameDataStreamReader定义于 readwriter.py L44内部持有SparkSession.readStream()返回的 Java 代理对象_jreader所有方法均通过 Py4J 转发到 JVM 端实现因此其在 Spark Classic 与 Spark Connect 模式下行为一致。2.1 通用配置format / schema / option / optionsspark.readStream.format(json) # 指定数据源格式如 json、parquet、csv、text spark.readStream.schema(col0 INT, col1 DOUBLE) # 指定 schemaStructType 或 DDL 字符串 spark.readStream.option(k, v) # 添加单个输入选项 spark.readStream.options(k1v1, k22) # 批量添加输入选项format(source)指定数据源名称。注意load()也可以直接接收format参数二者等价。schema(schema)接收pyspark.sql.types.StructType或 DDL 格式字符串如col0 INT, col1 DOUBLE。显式指定 schema 后底层数据源可跳过 schema 推断从而加速数据加载——这对 JSON 这类需要扫描数据才能推断类型的源尤其重要。option(key, value)/options(**options)向底层数据源传入键值选项值为None时会被静默跳过其余值统一经to_str转成字符串后传给 JVM 端_jreader.option(...)。例如 Rate 源可用option(rowsPerSecond, 10)控制每秒生成 10 行。2.2 文件类数据源json / csv / orc / parquet / text这五个方法分别加载对应格式的文件流均要求path为字符串传入非字符串会抛出PySparkTypeError。下表汇总各方法的核心专用参数均来自源码签名方法版本关键参数json(path, ...)2.0.0primitivesAsString、prefersDecimal、allowComments、allowUnquotedFieldNames、allowSingleQuotes、allowNumericLeadingZero、allowBackslashEscapingAnyCharacter、mode、columnNameOfCorruptRecord、dateFormat、timestampFormat、multiLine、allowUnquotedControlChars、lineSep、locale、dropFieldIfAllNull、encoding、pathGlobFilter、recursiveFileLookup、allowNonNumericNumbers、useUnsafeRowcsv(path, ...)2.0.0sep、encoding、quote、escape、comment、header、inferSchema、ignoreLeadingWhiteSpace、ignoreTrailingWhiteSpace、nullValue、nanValue、positiveInf、negativeInf、dateFormat、timestampFormat、maxColumns、maxCharsPerColumn、maxMalformedLogPerPartition、mode、columnNameOfCorruptRecord、multiLine、charToEscapeQuoteEscaping、enforceSchema、emptyValue、locale、lineSep、pathGlobFilter、recursiveFileLookup、unescapedQuoteHandlingorc(path, ...)2.3.0mergeSchema、pathGlobFilter、recursiveFileLookupparquet(path, ...)2.0.0mergeSchema、pathGlobFilter、recursiveFileLookup、datetimeRebaseMode、int96RebaseModetext(path, ...)2.0.0wholetext默认 False、lineSep、pathGlobFilter、recursiveFileLookup几个值得注意的语义细节JSON默认按 JSON Lines每行一条记录解析若一个文件包含一条 JSON 记录需设置multiLinetrue。未指定schema时会先扫描一遍输入以推断 schema。CSV开启inferSchema会先遍历一次数据确定 schema为避免全量扫描应关闭推断或显式传schema。text返回的 DataFrame schema 以名为value的字符串列开头后面可跟随分区列文件必须为 UTF-8 编码默认每行文本对应一行结果。ORC/Parquet共享pathGlobFilter文件通配过滤如*.json与recursiveFileLookup递归查找子目录两个文件发现选项Parquet 额外支持datetimeRebaseMode/int96RebaseMode控制日期时间与 INT96 的 rebase 行为。官方文档的 doctest 给出完整可运行范式——先写入一批静态文件再启动流式查询读取并输出到 console运行 3 秒后停止import tempfile, time with tempfile.TemporaryDirectory(prefixjson) as d: spark.createDataFrame( [(100, Hyukjin Kwon)], [age, name] ).write.mode(overwrite).format(json).save(d) q spark.readStream.schema(age INT, name STRING).json(d) \ .writeStream.format(console).start() time.sleep(3) q.stop()2.3 通用加载与表load / table / changes / nameload(pathNone, formatNone, schemaNone, **options)最通用的加载入口。format缺省时按spark.sql.sources.default配置默认 parquetpath非空字符串才合法为空会抛出VALUE_NOT_NON_EMPTY_STR。例如spark.readStream.format(rate).load()即使用内置 Rate 源持续生成timestamp、value两列数据是官方示例中最常用的测试源。table(tableName)3.1.0 引入在表上定义流式 DataFrame要求该表对应的 DataSource 支持流式模式。changes(tableName)4.2.0 引入返回指定表的**行级变更Change Data Capture**作为流式 DataFrame目前仅支持实现了TableCatalog.loadChangelog()的 Data Source V2 表可用option指定起始版本/时间戳spark.readStream.option(startingVersion, 10).changes(my_table)name(source_name)4.2.0 引入实验性为流式数据源命名用于在 checkpoint 元数据中标识该源从而为源的演进提供稳定的 checkpoint 位置。命名只允许 ASCII 字母、数字和下划线^[a-zA-Z0-9_]$违反时抛出INVALID_STREAMING_SOURCE_NAMEspark.readStream.format(rate).name(my_source)三、DataStreamWriter把流式 DataFrame 写到 SinkDataStreamWriter定义于 readwriter.py L996内部持有df._jdf.writeStream()的 Java 代理。它采用链式构建模式所有配置方法返回self最后由start()或toTable()真正启动查询并返回StreamingQuery对象。3.1 输出模式 outputModedf.writeStream.outputMode(append) # 仅写出新增行 df.writeStream.outputMode(complete) # 每次更新写出全部结果行 df.writeStream.outputMode(update) # 仅写出被更新的行append只把流式 DataFrame 中的新行写到 Sinkcomplete每次有更新时把全部结果行写到 Sink常用于全局聚合如df.groupby().count()update每次更新只写被更新的行若查询不含聚合等价于append。空字符串或非字符串参数会抛出VALUE_NOT_NON_EMPTY_STR。start()和toTable()也都接收outputMode参数效果与链式调用一致。3.2 写出配置format / option / options / partitionBy / queryNamedf.writeStream.format(console) # 指定 Sink 格式如 console、memory、parquet、csv df.writeStream.option(numRows, 3) # 单个输出选项console Sink 每批打印 3 行 df.writeStream.options(numRows3, truncateFalse) # 批量输出选项 df.writeStream.partitionBy(timestamp) # 按列在文件系统上分区Hive 风格 df.writeStream.queryName(streaming_query) # 为查询命名需在 Session 内唯一partitionBy支持单列或多列输出布局类似 Hive 分区方案。queryName必须在当前 SparkSession 的所有活跃查询中唯一空串会抛异常启动后可通过StreamingQuery.name读取。文件类 Sink如 parquet/csv写流通常必须提供checkpointLocation选项memory Sink 除外用于故障恢复。3.3 触发器 triggertrigger是控制流式查询执行节奏的核心参数5 种触发方式只能指定其一否则抛ONLY_ALLOW_SINGLE_TRIGGER不设置时默认尽可能快地执行等价于processingTime0 secondsdf.writeStream.trigger(processingTime5 seconds) # 周期性微批 df.writeStream.trigger(onceTrue) # 仅处理一批后终止 df.writeStream.trigger(continuous5 seconds) # 连续处理模式给定 checkpoint 间隔 df.writeStream.trigger(availableNowTrue) # 分批处理完所有可用数据后终止 df.writeStream.trigger(realTime5 seconds) # 实时模式按指定时长分批processingTime以处理时间为周期运行微批查询间隔为字符串如5 seconds、1 minuteonce处理一批数据后查询即终止True以外的值会抛VALUE_NOT_ALLOWEDcontinuous以给定 checkpoint 间隔运行连续处理查询实验性availableNow用多个批次处理完当前所有可用数据后自动终止适合有界数据的批式回放realTime按指定 batch 时长运行实时模式查询。底层实现上PySpark 通过_sc._jvm反射调用org.apache.spark.sql.streaming.Trigger的ProcessingTime/Once/Continuous/AvailableNow/RealTime工厂方法见 readwriter.py L1433-L1491。3.4 启动查询start / toTableq df.writeStream.trigger(processingTime5 seconds).start( queryNamethat_query, outputModeappend, formatmemory) q.name # that_query q.isActive # True q.stop()start(pathNone, formatNone, outputModeNone, partitionByNone, queryNameNone, **options)所有命名参数与链式调用方法等价适合在start处一次性集中配置path用于文件系统类 Sink通过**options可传入checkpointLocation等选项返回StreamingQuery可用来stop()、查看lastProgress等。toTable(tableName, ...)3.1.0 引入把流持续写入给定表。需要注意 v1/v2 表的差异见 readwriter.py L1890-L1897v1 表无论表是否存在都尊重partitionBy指定的分区列表不存在时会自动建表v2 表表已存在时忽略partitionBy仅当表不存在时生效且该 API 创建的 v2 表缺少自定义 properties、options、serde 等信息若需要这些能力应先手动创建 v2 表再执行。spark.readStream.format(rate).option(rowsPerSecond, 10).load() \ .writeStream.toTable(my_table2, queryNamethat_query, outputModeappend, formatparquet, checkpointLocation/path/to/checkpoint)3.5 任意下游处理foreach 与 foreachBatch当内置 Sink 无法满足需求时两个方法支持把流输出交给任意处理逻辑foreach(f)2.4.0 引入支持两种写法普通函数对每行调用f(row)。实现上被包装成func_without_process逐行迭代调用readwriter.py L1503-L1508。简单直接但无法在失败重放时对重复生成的数据去重。带process方法的对象推荐可实现精确一次语义class RowPrinter: def open(self, partition_id, epoch_id): print(Opened %d, %d % (partition_id, epoch_id)) return True def process(self, row): print(row) def close(self, error): print(Closed with error: %s % str(error)) q df.writeStream.foreach(RowPrinter()).start()对象的生命周期约定详见 readwriter.py L1594-L1648open(partition_id, epoch_id)可选初始化处理资源打开连接、启动事务等返回True才继续处理本批数据process(row)必选处理每一行close(error)可选全部行处理完毕后收尾关闭连接、提交事务出错时收到 error 对象每个任务对应对象的一个副本负责该任务某个分区的全部数据对象必须可序列化每个任务会拿到反序列化后的全新副本因此连接、事务等初始化应在open()之后进行在微批模式下(partition_id, epoch_id)元组唯一对应同一批数据可用于去重与事务性提交实现 exactly-once但在连续处理模式下该保证不成立不能用于去重。foreachBatch(func)2.4.0 引入每个微批调用一次func(batch_df, batch_id)其中batch_df是本批输出 DataFramebatch_id可用于去重和事务性写入——同一batchId对应完全相同的输出数据假设查询内操作确定。限制仅支持微批执行模式trigger 不能是 continuous在Spark Connect模式下函数无法访问外部定义的变量经典模式下可以二者行为不同见 readwriter.py L1720-L1740。四、从源码看 I/O 链路的实现与验证参数转发机制DataStreamReader/DataStreamWriter均持有 Java 代理_jreader/_jwritePython 端只做类型校验如PySparkTypeError/PySparkValueError与字符串化随后直接调用 JVM 端方法。OptionUtils._set_opts负责把大量具名参数批量映射为底层 option。foreach 的 JVM 桥接foreach通过CPickleSerializer序列化 Python 函数再包装为org.apache.spark.sql.execution.python.streaming.PythonForeachWriter交给 JVM 端流式执行readwriter.py L1690-L1702foreachBatch则通过PythonForeachBatchHelper.callForeachBatch注册回调服务器readwriter.py L1747-L1751。checkpoint 与容错start()/toTable()的**options中通常需要checkpointLocationmemory Sink 除外配合DataStreamReader.name()的源命名即可实现稳定的 checkpoint 元数据与源演进。测试验证每个方法都附有 doctest 示例_test()以local[4]启动 SparkSession 执行全量 doctestreadwriter.py L1936-L1960结合 python/docs/source/reference/pyspark.ss/query_management.rst 中StreamingQuery的lastProgress/processAllAvailable/awaitTermination等方法可完成启动 → 观察 → 停止的完整查询生命周期管理。五、端到端实战示例把前三节串起来一个完整的结构化流任务通常遵循readStream 配置源 → 流式变换 → writeStream 配置 Sink/触发 → start的链路import time from pyspark.sql import SparkSession spark SparkSession.builder.master(local[2]).appName(streaming-io-demo).getOrCreate() # 1) 读取Rate 源每秒 10 行并显式命名源 df spark.readStream.format(rate) \ .option(rowsPerSecond, 10) \ .name(rate_source) \ .load() # 2) 变换按 3 取模 df df.selectExpr(value % 3 as v) # 3) 写出complete 输出模式 周期微批写到 memory 表 q df.writeStream \ .outputMode(complete) \ .queryName(mod_counts) \ .trigger(processingTime2 seconds) \ .format(memory) \ .start() time.sleep(5) spark.sql(SELECT * FROM mod_counts).show() # 用普通 SQL 查看聚合结果 q.stop()如果目标是文件类 Sink则使用start指定输出目录并携带 checkpointq df.writeStream.format(parquet) \ .outputMode(append) \ .option(checkpointLocation, /tmp/ckpt) \ .partitionBy(v) \ .start(/tmp/stream_out)需要接入任意外部系统数据库、消息队列、自定义存储时改用foreach对象或foreachBatch函数即可无缝扩展。结语PySpark Structured Streaming 的 I/O 层以DataStreamReader13 个方法与DataStreamWriter11 个方法两个对称接口覆盖了流式数据读写的完整生命周期从文件、表、变更日志到任意自定义源从内置 Sink 到foreach/foreachBatch的自定义写出再到 5 种触发器与 3 种输出模式的灵活组合。理解本文梳理的每个方法的参数语义、版本引入时间table/toTable3.1.0、foreach/foreachBatch2.4.0、changes/name4.2.0以及 readwriter.py 中对应的实现细节将帮助你在实际项目中正确、稳定地构建端到端流处理管道。赞分享大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载相关推荐Data Engineering Zoomcamp 实战用 PySpark Structured Streaming 实时消费 Kafka 出租车数据Data Engineering Zoomcamp 实战用 PySpark Structured Streaming 实时消费 Kafka 出租车数据 本篇技教程数据工程Data Engineering Zoomcamp 实战使用 Redpanda 运行 PySpark Structured Streaming 流式管道Data Engineering Zoomcamp 实战使用 Redpanda 运行 PySpark Structured Streaming 流式管道 本指教程数据工程Data Engineering Zoomcamp 实战使用 PySpark Structured Streaming 消费 Kafka 流并完成窗口聚合Data Engineering Zoomcamp 实战使用 PySpark Structured Streaming 消费 Kafka 流并完成窗口聚合 本教程数据工程上一篇NeRF与SLAM融合从iMAP到NICE-SLAM的终极演进指南下一篇领域事件序列化Modular Monolith DDD JSON与多态处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考