从 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 文档贡献者;本文为授权中文整理。

版本和运行环境
检索日为 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.












暂无评论内容