为 PySpark 注册按行和 Arrow 批次读取的数据源

从 Spark 4.0 开始,Python Data Source API 允许开发者用 Python 定义自己的数据源和写入端,并通过熟悉的 DataFrame 读写接口调用。一个读取型数据源最核心的分工是:DataSource 提供名称、模式与读取器;DataSourceReader 为分区产生数据;Spark 把产出转换成 DataFrame。

本文依据 Apache Spark 官方 Python Data Source API 全文核查后,完整整理 SimpleDataSource 与 ArrowBatchDataSource 两条批读取主线。原页的 Faker、流式读写及准入控制示例不作为本文可用链路;有关缺项在文末说明。原作与维护方:Apache Software Foundation 及 Spark 文档贡献者;本文为授权中文整理。

两种PySpark自定义数据源都先注册名称与schema,再由Reader按分区产出数据;上路逐行yield元组,下路yield列式Arrow RecordBatch,最后转成DataFrame。
原创技术示意图:同一个 Data Source 接口,两种数据产出粒度;未完纪编辑整理。

版本和运行环境

检索日为 2026-10-05,latest 页面标题标识 PySpark 4.2.0。它是移动链接,不代表每个 Spark 4.0 及后续发行版都具有同样的直接 Arrow 批次功能与自动注册机制。实施时应选定 Spark/PySpark 的一致版本,并查看该版本的 API 文档。

Simple 示例无需 Faker 等业务数据生成库,但仍需要正常的 PySpark 运行环境。Arrow 示例直接导入 PyArrow,还需确保它在执行读取器的 Python worker 环境中可用。检索时官方安装页列出 Java 17 或以上,Spark SQL 依赖包括 pandas 2.2.0 至 3.0.0 之前、PyArrow 18.0.0 或以上;这些是所读取版本的要求,不应无条件套用到旧版集群。

官方提供 pip install "pyspark[sql]" 这一依赖组合入口。实际复现应在独立虚拟环境中固定选定版本,并保证驱动端、执行端和 Python 依赖一致。本稿没有安装依赖、启动 Spark 或连接现有集群;所有输出说明来自代码结构或原文示例。

第一条主线:定义一个产生两行数据的源

先继承 DataSource,给它一个不与内建或 JVM 数据源冲突的短名称。schema() 返回包含 name 字符串和 age 整数的 StructType。reader() 接收 Spark 传入的模式并返回读取器;这个最小示例固定两列,因此不使用该参数来动态改变输出。

from typing import Iterator, Tuple

from pyspark.sql.datasource import DataSource, DataSourceReader, InputPartition
from pyspark.sql.types import IntegerType, StringType, StructField, StructType


class SimpleDataSource(DataSource):
    """为 PySpark 生成恰好两行合成数据。"""

    @classmethod
    def name(cls) -> str:
        return "simple"

    def schema(self) -> StructType:
        return StructType([
            StructField("name", StringType()),
            StructField("age", IntegerType())
        ])

    def reader(self, schema: StructType) -> DataSourceReader:
        return SimpleDataSourceReader()


class SimpleDataSourceReader(DataSourceReader):

    def read(self, partition: InputPartition) -> Iterator[Tuple]:
        yield ("Alice", 20)
        yield ("Bob", 30)

这里的 yield 每次交出一个元组,元组列数、位置与类型必须符合模式。InputPartition 用来表示当前分区的读取输入;这段最小读取器没有使用它,也没有扩展分区规划。它演示接口连接方式,不是可伸缩的数据分片策略。

注册,再用 format 读取

from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()
spark.dataSource.register(SimpleDataSource)

spark.read.format("simple").load().show()

注册传入的是类 SimpleDataSource,不是手动创建的实例。format("simple") 对应 name() 的返回值。原文展示的输出是两列两行:

+-----+---+
| name|age|
+-----+---+
|Alice| 20|
|  Bob| 30|
+-----+---+

这是原文预期输出,不是本次测试结果。真实 DataFrame 若需要稳定排序,应显式排序;不能从入门展示推导分布式系统中的普遍行序保证。

需要注意,调用方随意用 .schema(...) 指定其他列,并不会使这个读取器自动改写固定元组。若要支持用户提供的模式,应在读取器中验证并匹配它,或明确拒绝不兼容模式;不能只改变声明而保持产出不变。

第二条主线:直接产生 Arrow RecordBatch

逐行构造 Python 对象会增加数据交换开销。支持直接 Arrow 批次的读取器可以在 read() 中产出 pyarrow.RecordBatch,把多行数据按列式批次交给 Spark。原文称大型数据集可能获得最高一个数量级的改善;这是原文的概括,本稿未复现基准,也不把它当成所有负载的保证。批大小、类型转换、数据源开销、并行度和内存占用都会影响结果。

下面保留原文完整数据与单分区结构,并作两处静态整理:把 reader() 的模式类型标注从原文的 str 改为 StructType;按照同页的序列化建议,把 PyArrow 导入放进使用它的 read()。改动没有经过运行测试。

from pyspark.sql import SparkSession
from pyspark.sql.datasource import DataSource, DataSourceReader, InputPartition
from pyspark.sql.types import StructType


class ArrowBatchDataSource(DataSource):
    """演示直接返回 Arrow RecordBatch 的 Python 数据源。"""

    @classmethod
    def name(cls):
        return "arrowbatch"

    def schema(self):
        return "key int, value string"

    def reader(self, schema: StructType):
        return ArrowBatchDataSourceReader(schema, self.options)


class ArrowBatchDataSourceReader(DataSourceReader):

    def __init__(self, schema: StructType, options):
        self.schema = schema
        self.options = options

    def read(self, partition):
        import pyarrow as pa

        keys = pa.array([1, 2, 3, 4, 5], type=pa.int32())
        values = pa.array(
            ["one", "two", "three", "four", "five"],
            type=pa.string()
        )
        schema = pa.schema([
            ("key", pa.int32()),
            ("value", pa.string())
        ])
        record_batch = pa.RecordBatch.from_arrays(
            [keys, values], schema=schema
        )
        yield record_batch

    def partitions(self):
        num_part = 1
        return [InputPartition(i) for i in range(num_part)]


spark = SparkSession.builder.appName("ArrowBatchExample").getOrCreate()
spark.dataSource.register(ArrowBatchDataSource)

df = spark.read.format("arrowbatch").load()
df.show()

schema() 使用 DDL 字符串声明 key int, value string;传到读取器中的已解析模式则是 StructType。Arrow 侧显式声明 int32 和 string,与 Spark 的整数和字符串列相对应。两个 Arrow 数组长度同为五,组成一个五行两列的 RecordBatch,再通过一次 yield 交给引擎。

阅读代码可以确定合成数据对为 1/one 至 5/five,但这里不伪造运行终端输出。与前例一样,这个读取器保存了传入的模式与选项,却仍然生成固定两列;它没有实现任意 schema、投影或外部数据读取。

一个分区不等于可以随意增加分区数

partitions() 返回含一个 InputPartition 的列表。示例中的 read() 不查看分区值,因此若只把 num_part 改大,每个分区都会生成同样五条记录,造成重复。真正的多分区数据源应把分片范围、文件段或其他可序列化的定位信息放入分区对象,并由 read() 只读取自己的部分。

同样,不能把整个外部数据集一次装进一个 RecordBatch 就认为完成了批处理。输入很大时,需要有界地构建批次并逐批产出,监控 Arrow 缓冲区与 Python worker 内存;具体批大小必须依据负载测量。

序列化、注册与名称解析

官方要求用户定义的 DataSource、Reader、Writer、StreamReader、StreamWriter 及其方法能被 pickle 序列化。方法中使用的运行期库应在方法内部导入;原文用 TaskContext 展示这一原则。不要把打开的网络连接、不可序列化的句柄、包含秘密的运行对象当成读取器实例状态随意传往 worker。

这并不意味着可以反序列化来源不明的 pickle。数据源注册的是可执行 Python 代码,应当来自受信任的项目及包。本文两个固定数据示例没有 SQL 拼接、shell 命令、动态执行外部输入或硬编码凭据;这只说明本次示例中未见这些入口,并不构成数据源框架安全保证。

名称解析还有三个需要保留的规则:

  • 同名时,内建数据源以及 Scala/Java 数据源优先于 Python 数据源,所以 simple、arrowbatch 之类名称也应在实际会话中检查冲突。
  • 可以重复注册同名 Python 数据源;后注册的会覆盖先注册的 Python 版本。长生命周期会话中不要把同名重注册当成无影响操作。
  • 当前源页还说明,顶层模块名以 pyspark_ 开头,并导出 DefaultSource 时,可以自动注册。原文以 pyspark_huggingface 为参考。这个能力应按所选版本核对,本文没有安装或检验该第三方包。

为什么没有把原页所有流式代码拼进来

Python Data Source API 的能力不限于本文的批读取。原文列出的对应关系是:批读取实现 reader(),批写入实现 writer();流式读取实现 streamReader() 或 simpleStreamReader(),流式写入实现 streamWriter()。较低吞吐且不需要分区的场景可考虑简单流式读取器;若同时实现两个读取入口,原文说明 streamReader() 优先。

不过,原页的综合示例不能不加审查地串成一个完整应用:FakeDataSource.schema() 声明四列,而普通流式读取器产生两列,简单流式读取器产生一列;流式 SimpleCommitMessage 只有属性注解,缺少批写入示例中使用的 @dataclass 或等价构造器,却用关键字实参创建;流式提交代码还使用 os 和 json 而相邻片段没有相应导入。它的本地文件写入、错误提交和 None 消息处理也需要独立设计。

同页的准入控制例子通过 getDefaultReadLimit() 返回 ReadMaxRows(20),并要求 latestOffset(start, limit) 尊重引擎传入的限制,但其完整数据源、状态和可用记录上界并没有在片段中给齐。它与偏移量重放、失败恢复、幂等提交是另一个完整主题。本文没有补写这些缺口再宣称完成流式教程。

核查与许可

两条主线保留了注册、模式、读取器、分区与调用流程;新增的是类型注解订正、方法内导入、分区重复与模式一致性提醒。未执行任何 Spark 作业,未进行速度、内存或分布式兼容性测试,文中预期结果不构成测试证据。

原文版权归 The Apache Software Foundation,页面声明采用 Apache License 2.0;本文保留来源和许可链接,并明确标识翻译与修改。Apache Spark、PySpark 等名称的使用不表示官方背书。

原始版权与许可

Apache Spark;Copyright 2014 and onwards The Apache Software Foundation. This product includes software developed at The Apache Software Foundation (http://www.apache.org/). 此处保留与本文引用内容相关的原始 NOTICE;官方 NOTICE 来源。其他组件的密码学与Metrics说明不属于本文转载内容。

                                 Apache License
                           Version 2.0, January 2004
                        http://www.apache.org/licenses/

   TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION

   1. Definitions.

      "License" shall mean the terms and conditions for use, reproduction,
      and distribution as defined by Sections 1 through 9 of this document.

      "Licensor" shall mean the copyright owner or entity authorized by
      the copyright owner that is granting the License.

      "Legal Entity" shall mean the union of the acting entity and all
      other entities that control, are controlled by, or are under common
      control with that entity. For the purposes of this definition,
      "control" means (i) the power, direct or indirect, to cause the
      direction or management of such entity, whether by contract or
      otherwise, or (ii) ownership of fifty percent (50%) or more of the
      outstanding shares, or (iii) beneficial ownership of such entity.

      "You" (or "Your") shall mean an individual or Legal Entity
      exercising permissions granted by this License.

      "Source" form shall mean the preferred form for making modifications,
      including but not limited to software source code, documentation
      source, and configuration files.

      "Object" form shall mean any form resulting from mechanical
      transformation or translation of a Source form, including but
      not limited to compiled object code, generated documentation,
      and conversions to other media types.

      "Work" shall mean the work of authorship, whether in Source or
      Object form, made available under the License, as indicated by a
      copyright notice that is included in or attached to the work
      (an example is provided in the Appendix below).

      "Derivative Works" shall mean any work, whether in Source or Object
      form, that is based on (or derived from) the Work and for which the
      editorial revisions, annotations, elaborations, or other modifications
      represent, as a whole, an original work of authorship. For the purposes
      of this License, Derivative Works shall not include works that remain
      separable from, or merely link (or bind by name) to the interfaces of,
      the Work and Derivative Works thereof.

      "Contribution" shall mean any work of authorship, including
      the original version of the Work and any modifications or additions
      to that Work or Derivative Works thereof, that is intentionally
      submitted to Licensor for inclusion in the Work by the copyright owner
      or by an individual or Legal Entity authorized to submit on behalf of
      the copyright owner. For the purposes of this definition, "submitted"
      means any form of electronic, verbal, or written communication sent
      to the Licensor or its representatives, including but not limited to
      communication on electronic mailing lists, source code control systems,
      and issue tracking systems that are managed by, or on behalf of, the
      Licensor for the purpose of discussing and improving the Work, but
      excluding communication that is conspicuously marked or otherwise
      designated in writing by the copyright owner as "Not a Contribution."

      "Contributor" shall mean Licensor and any individual or Legal Entity
      on behalf of whom a Contribution has been received by Licensor and
      subsequently incorporated within the Work.

   2. Grant of Copyright License. Subject to the terms and conditions of
      this License, each Contributor hereby grants to You a perpetual,
      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
      copyright license to reproduce, prepare Derivative Works of,
      publicly display, publicly perform, sublicense, and distribute the
      Work and such Derivative Works in Source or Object form.

   3. Grant of Patent License. Subject to the terms and conditions of
      this License, each Contributor hereby grants to You a perpetual,
      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
      (except as stated in this section) patent license to make, have made,
      use, offer to sell, sell, import, and otherwise transfer the Work,
      where such license applies only to those patent claims licensable
      by such Contributor that are necessarily infringed by their
      Contribution(s) alone or by combination of their Contribution(s)
      with the Work to which such Contribution(s) was submitted. If You
      institute patent litigation against any entity (including a
      cross-claim or counterclaim in a lawsuit) alleging that the Work
      or a Contribution incorporated within the Work constitutes direct
      or contributory patent infringement, then any patent licenses
      granted to You under this License for that Work shall terminate
      as of the date such litigation is filed.

   4. Redistribution. You may reproduce and distribute copies of the
      Work or Derivative Works thereof in any medium, with or without
      modifications, and in Source or Object form, provided that You
      meet the following conditions:

      (a) You must give any other recipients of the Work or
          Derivative Works a copy of this License; and

      (b) You must cause any modified files to carry prominent notices
          stating that You changed the files; and

      (c) You must retain, in the Source form of any Derivative Works
          that You distribute, all copyright, patent, trademark, and
          attribution notices from the Source form of the Work,
          excluding those notices that do not pertain to any part of
          the Derivative Works; and

      (d) If the Work includes a "NOTICE" text file as part of its
          distribution, then any Derivative Works that You distribute must
          include a readable copy of the attribution notices contained
          within such NOTICE file, excluding those notices that do not
          pertain to any part of the Derivative Works, in at least one
          of the following places: within a NOTICE text file distributed
          as part of the Derivative Works; within the Source form or
          documentation, if provided along with the Derivative Works; or,
          within a display generated by the Derivative Works, if and
          wherever such third-party notices normally appear. The contents
          of the NOTICE file are for informational purposes only and
          do not modify the License. You may add Your own attribution
          notices within Derivative Works that You distribute, alongside
          or as an addendum to the NOTICE text from the Work, provided
          that such additional attribution notices cannot be construed
          as modifying the License.

      You may add Your own copyright statement to Your modifications and
      may provide additional or different license terms and conditions
      for use, reproduction, or distribution of Your modifications, or
      for any such Derivative Works as a whole, provided Your use,
      reproduction, and distribution of the Work otherwise complies with
      the conditions stated in this License.

   5. Submission of Contributions. Unless You explicitly state otherwise,
      any Contribution intentionally submitted for inclusion in the Work
      by You to the Licensor shall be under the terms and conditions of
      this License, without any additional terms or conditions.
      Notwithstanding the above, nothing herein shall supersede or modify
      the terms of any separate license agreement you may have executed
      with Licensor regarding such Contributions.

   6. Trademarks. This License does not grant permission to use the trade
      names, trademarks, service marks, or product names of the Licensor,
      except as required for reasonable and customary use in describing the
      origin of the Work and reproducing the content of the NOTICE file.

   7. Disclaimer of Warranty. Unless required by applicable law or
      agreed to in writing, Licensor provides the Work (and each
      Contributor provides its Contributions) on an "AS IS" BASIS,
      WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
      implied, including, without limitation, any warranties or conditions
      of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
      PARTICULAR PURPOSE. You are solely responsible for determining the
      appropriateness of using or redistributing the Work and assume any
      risks associated with Your exercise of permissions under this License.

   8. Limitation of Liability. In no event and under no legal theory,
      whether in tort (including negligence), contract, or otherwise,
      unless required by applicable law (such as deliberate and grossly
      negligent acts) or agreed to in writing, shall any Contributor be
      liable to You for damages, including any direct, indirect, special,
      incidental, or consequential damages of any character arising as a
      result of this License or out of the use or inability to use the
      Work (including but not limited to damages for loss of goodwill,
      work stoppage, computer failure or malfunction, or any and all
      other commercial damages or losses), even if such Contributor
      has been advised of the possibility of such damages.

   9. Accepting Warranty or Additional Liability. While redistributing
      the Work or Derivative Works thereof, You may choose to offer,
      and charge a fee for, acceptance of support, warranty, indemnity,
      or other liability obligations and/or rights consistent with this
      License. However, in accepting such obligations, You may act only
      on Your own behalf and on Your sole responsibility, not on behalf
      of any other Contributor, and only if You agree to indemnify,
      defend, and hold each Contributor harmless for any liability
      incurred by, or claims asserted against, such Contributor by reason
      of your accepting any such warranty or additional liability.

   END OF TERMS AND CONDITIONS

   APPENDIX: How to apply the Apache License to your work.

      To apply the Apache License to your work, attach the following
      boilerplate notice, with the fields enclosed by brackets "[]"
      replaced with your own identifying information. (Don't include
      the brackets!)  The text should be enclosed in the appropriate
      comment syntax for the file format. We also recommend that a
      file or class name and description of purpose be included on the
      same "printed page" as the copyright notice for easier
      identification within third-party archives.

   Copyright [yyyy] [name of copyright owner]

   Licensed under the Apache License, Version 2.0 (the "License");
   you may not use this file except in compliance with the License.
   You may obtain a copy of the License at

       http://www.apache.org/licenses/LICENSE-2.0

   Unless required by applicable law or agreed to in writing, software
   distributed under the License is distributed on an "AS IS" BASIS,
   WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
   See the License for the specific language governing permissions and
   limitations under the License.
© 版权声明
THE END
喜欢就支持一下吧
点赞0 分享
评论 抢沙发

请登录后发表评论

    暂无评论内容