- 大数据
- 数据分析
- 批处理
- 流处理
- 机器学习
- 图计算
【免费下载链接】spark
Apache Spark - A unified analytics engine for large-scale data processing
导读
本文聚焦 Apache Spark 仓库中python/pyspark/sql/tests/coercion/pandas_2/golden_python_udf_input_type_coercion_with_arrow_and_pandas.md这份 golden(黄金基线)测试数据,逐行解读在Arrow 优化 + Legacy Pandas 转换双重开启的情况下,Python UDF 的输入参数从 Spark SQL 类型强转为 Python 对象的完整行为矩阵。你将掌握:39 个覆盖标量、嵌套复杂类型与空值场景的输入强转预期结果、该测试与纯 Arrow 路径、纯 Pickle(vanilla)路径的差异点,以及 pandas 2 / pandas 3 主版本差异如何影响 UDF 输入的实际类型,从而在编写与排查 PySpark Python UDF 时准确预判传入 Python 函数的值类型。
一、背景:为什么需要一份“输入类型强转”的基线测试
在 PySpark 中,Python UDF 的输入参数需要从 JVM 侧传入的 Spark SQL 类型转换成 Python 侧的实际对象。Spark 默认走Pickle(pickle-based batch)序列化;当开启 Arrow 优化后,则改走Arrow 列式内存 + pandas 转换的路径。两条路径的类型强转规则并不一致,同一列数据传入 UDF 后可能得到不同的 Python 类型(例如int与numpy.int32、datetime与pandas.Timestamp、None与float('nan'))。
为了让这一行为可预期、可回归,Spark 的测试套件以golden 文件形式固化每一组“Spark 类型 + Spark 输入值”在某种执行路径下应当得到的“Python 类型 + Python 值”,并用自动化断言防止任何一条路径的类型强转行为发生未登记的变化。本文讨论的with_arrow_and_pandas版本正是其中输入侧最复杂的一条路径:Arrow 负责列式传输,Legacy pandas 转换负责将 Arrow 数组转换为 UDF 函数实际接收到的 Python 对象。
二、三条执行路径与两个配置开关
在 test_python_udf_input_type.py 中,同一个测试用例集被跑在三种配置组合下:
| 测试方法 | spark.sql.execution.pythonUDF.arrow.enabled | spark.sql.legacy.execution.pythonUDF.pandas.conversion.enabled | golden 文件 |
|---|---|---|---|
test_python_input_type_coercion_vanilla | False | False | golden_python_udf_input_type_coercion_vanilla |
test_python_input_type_coercion_with_arrow | True | False | golden_python_udf_input_type_coercion_with_arrow |
test_python_input_type_coercion_with_arrow_and_pandas | True | True | pandas_2/golden_python_udf_input_type_coercion_with_arrow_and_pandas(或pandas_3/...) |
三个测试方法最终都调用_run_udf_input_type_coercion,其核心代码如下:
def _run_udf_input_type_coercion(self, use_arrow, legacy_pandas, golden_file, test_name): with self.sql_conf( { "spark.sql.execution.pythonUDF.arrow.enabled": use_arrow, "spark.sql.legacy.execution.pythonUDF.pandas.conversion.enabled": legacy_pandas, } ): self._compare_or_generate_golden(golden_file, test_name)需要说明的是:Arrow 优化在 PySpark 中默认关闭。在 udf.py 的_create_py_udf中,仅当配置spark.sql.execution.pythonUDF.arrow.enabled == "true"(或通过useArrow显式指定)且环境中装有 pandas、pyarrow 时,UDF 的求值类型才会从SQL_BATCHED_UDF提升为SQL_ARROW_BATCHED_UDF;缺少依赖时打印警告并回退到普通路径。这解释了为什么类型强转的差异需要单独用一份 golden 文件来追踪——它只会在用户显式开启 Arrow 优化后出现。
而spark.sql.legacy.execution.pythonUDF.pandas.conversion.enabled这一开关控制的是 Worker 端将 Arrow 数据转换成 pandas 对象时是否走“legacy 转换路径”,其默认值为false,见 worker_util.py 中RunnerConf.use_legacy_pandas_udf_conversion的定义。本文的关联文档记录的就是该开关为true时的行为快照。
三、Golden 文件的生成与校验机制
Golden 测试的底座是 goldenutils.py 中的GoldenFileTestMixin,它负责三件关键的事:
- 时区固定:测试类在
setUpClass中通过setup_timezone将进程时区固定为America/Los_Angeles,并同步到 Spark 会话的spark.sql.session.timeZone,保证date/timestamp相关用例在任何机器上产出相同的字符串(teardown_timezone负责还原)。这解释了 golden 表中datetime.date、datetime.datetime值的确定性来源。 - 生成与比对双模式:读取环境变量
SPARK_GENERATE_GOLDEN_FILES,为"1"时执行save_golden写出 CSV 与 Markdown 两份基线;否则load_golden_csv读入既有基线并逐行比对。Markdown 的写出依赖tabulate包,未安装时仅警告而不中断。 - 类型字符串规范化:
repr_type对 SparkDataType使用simpleString()(如tinyint、array<int>、struct<a1:int,a2:string>),对 Python 类型使用__name__,从而形成 golden 表中“Spark Type”与“Python Type”两列的统一书写。
测试代码头部还给出了重新生成 golden 文件的标准命令:
SPARK_GENERATE_GOLDEN_FILES=1 python/run-tests -k \ --testnames 'pyspark.sql.tests.coercion.test_python_udf_input_type'四、测试用例的构造方式
在 test_python_udf_input_type.py 中,39 个用例以(case_name, spark_type, data_func)三元组形式注册。对每种类型,测试构造一个value列 DataFrame 并repartition(1),随后选出三路 UDF:
def type_udf(x): if x is None: return "NoneType" else: return type(x).__name__ def value_udf(x): return x def value_str(x): return str(x) type_test_udf = udf(type_udf, returnType=StringType()) value_test_udf = udf(value_udf, returnType=spark_type) value_str_udf = udf(value_str, returnType=StringType()) result_df = input_df.select( value_test_udf("value").alias("python_value"), type_test_udf("value").alias("python_type"), value_str_udf("value").alias("python_value_str"), )其中value_udf原样返回输入,并用assert values == input_data校验值本身能原样往返(即强转后数值不丢失);type_udf记录每个元素在 Python 侧的实际类型名;value_str_udf记录str()后的实际值字符串。若执行抛出异常(如某些空值断言失败),测试会把异常信息以✗ ...前缀写入 golden 单元格,而不是直接失败——这样连“转换失败”的行为本身也被登记为基线。
五、完整 Golden 基线数据(pandas 2)
以下为 pandas_2/golden_python_udf_input_type_coercion_with_arrow_and_pandas.md 的完整 39 行数据,该表同时对应同目录下的 CSV 版本:
| Test Case | Spark Type | Spark Value | Python Type | Python Value | |
|---|---|---|---|---|---|
| 0 | byte_values | tinyint | [-128, 127, 0] | ['int', 'int', 'int'] | ['-128', '127', '0'] |
| 1 | byte_null | tinyint | [None, 42] | ['float', 'float'] | ['nan', '42.0'] |
| 2 | short_values | smallint | [-32768, 32767, 0] | ['int', 'int', 'int'] | ['-32768', '32767', '0'] |
| 3 | short_null | smallint | [None, 123] | ['float', 'float'] | ['nan', '123.0'] |
| 4 | int_values | int | [-2147483648, 2147483647, 0] | ['int', 'int', 'int'] | ['-2147483648', '2147483647', '0'] |
| 5 | int_null | int | [None, 456] | ['float', 'float'] | ['nan', '456.0'] |
| 6 | long_values | bigint | [-9223372036854775808, 9223372036854775807, 0] | ['int', 'int', 'int'] | ['-9223372036854775808', '9223372036854775807', '0'] |
| 7 | long_null | bigint | [None, 789] | ['float', 'float'] | ['nan', '789.0'] |
| 8 | float_values | float | [0.0, 1.0, 3.140000104904175] | ['float', 'float', 'float'] | ['0.0', '1.0', '3.140000104904175'] |
| 9 | float_null | float | [None, 3.140000104904175] | ['float', 'float'] | ['nan', '3.140000104904175'] |
| 10 | double_values | double | [0.0, 1.0, 0.3333333333333333] | ['float', 'float', 'float'] | ['0.0', '1.0', '0.3333333333333333'] |
| 11 | double_null | double | [None, 2.71] | ['float', 'float'] | ['nan', '2.71'] |
| 12 | decimal_values | decimal(3,2) | [Decimal('5.35'), Decimal('1.23')] | ['Decimal', 'Decimal'] | ['5.35', '1.23'] |
| 13 | decimal_null | decimal(3,2) | [None, Decimal('9.99')] | ['NoneType', 'Decimal'] | ['None', '9.99'] |
| 14 | string_values | string | ['abc', '', 'hello'] | ['str', 'str', 'str'] | ['abc', '', 'hello'] |
| 15 | string_null | string | [None, 'test'] | ['NoneType', 'str'] | ['None', 'test'] |
| 16 | binary_values | binary | [b'abc', b'', b'ABC'] | ['bytes', 'bytes', 'bytes'] | ["b'abc'", "b''", "b'ABC'"] |
| 17 | binary_null | binary | [None, b'test'] | ['NoneType', 'bytes'] | ['None', "b'test'"] |
| 18 | boolean_values | boolean | [True, False] | ['bool', 'bool'] | ['True', 'False'] |
| 19 | boolean_null | boolean | [None, True] | ['NoneType', 'bool'] | ['None', 'True'] |
| 20 | date_values | date | [datetime.date(2020, 2, 2), datetime.date(1970, 1, 1)] | ['date', 'date'] | ['2020-02-02', '1970-01-01'] |
| 21 | date_null | date | [None, datetime.date(2023, 1, 1)] | ['NoneType', 'date'] | ['None', '2023-01-01'] |
| 22 | timestamp_values | timestamp | [datetime.datetime(2020, 2, 2, 12, 15, 16, 123000)] | ['Timestamp'] | ['2020-02-02 12:15:16.123000'] |
| 23 | timestamp_null | timestamp | [None, datetime.datetime(2023, 1, 1, 12, 0)] | ['NaTType', 'Timestamp'] | ['NaT', '2023-01-01 12:00:00'] |
| 24 | array_int_values | array<int> | [[1, 2, 3], [], [1, None, 3]] | ['list', 'list', 'list'] | ['[1, 2, 3]', '[]', '[1, None, 3]'] |
| 25 | array_int_null | array<int> | [None, [4, 5, 6]] | ['NoneType', 'list'] | ['None', '[np.int32(4), np.int32(5), np.int32(6)]'] |
| 26 | map_str_int_values | map<string,int> | [{'world': 2, 'hello': 1}, {}] | ['dict', 'dict'] | ["{'world': 2, 'hello': 1}", '{}'] |
| 27 | map_str_int_null | map<string,int> | [None, {'test': 123}] | ['NoneType', 'dict'] | ['None', "{'test': 123}"] |
| 28 | struct_int_str_values | struct<a1:int,a2:string> | [Row(a1=1, a2='hello'), Row(a1=2, a2='world')] | ['Row', 'Row'] | ["Row(a1=1, a2='hello')", "Row(a1=2, a2='world')"] |
| 29 | struct_int_str_null | struct<a1:int,a2:string> | [None, Row(a1=99, a2='test')] | ['NoneType', 'Row'] | ['None', "Row(a1=99, a2='test')"] |
| 30 | array_array_int | array<array<int>> | [[[1, 2, 3]], [[1], [2, 3]]] | ['list', 'list'] | ['[[np.int32(1), np.int32(2), np.int32(3)]]', '[[np.int32(1)], [np.int32(2), np.int32(3)]]'] |
| 31 | array_map_str_int | array<map<string,int>> | [[{'world': 2, 'hello': 1}], [{'a': 1}, {'b': 2}]] | ['list', 'list'] | ["[{'world': 2, 'hello': 1}]", "[{'a': 1}, {'b': 2}]"] |
| 32 | array_struct_int_str | array<struct<a1:int,a2:string>> | [[Row(a1=1, a2='hello')], [Row(a1=1, a2='hello'), Row(a1=2, a2='world')]] | ['list', 'list'] | ["[Row(a1=1, a2='hello')]", "[Row(a1=1, a2='hello'), Row(a1=2, a2='world')]"] |
| 33 | map_int_array_int | map<int,array<int>> | [{1: [1, 2, 3]}, {1: [1], 2: [2, 3]}] | ['dict', 'dict'] | ['{1: [np.int32(1), np.int32(2), np.int32(3)]}', '{1: [np.int32(1)], 2: [np.int32(2), np.int32(3)]}'] |
| 34 | map_int_map_str_int | map<int,map<string,int>> | [{1: {'world': 2, 'hello': 1}}] | ['dict'] | ["{1: {'world': 2, 'hello': 1}}"] |
| 35 | map_int_struct_int_str | map<int,struct<a1:int,a2:string>> | [{1: Row(a1=1, a2='hello')}] | ['dict'] | ["{1: Row(a1=1, a2='hello')}"] |
| 36 | struct_int_array_int | struct<a:int,b:array<int>> | [Row(a=1, b=[1, 2, 3])] | ['Row'] | ['Row(a=1, b=[np.int32(1), np.int32(2), np.int32(3)])'] |
| 37 | struct_int_map_str_int | struct<a:int,b:map<string,int>> | [Row(a=1, b={'world': 2, 'hello': 1})] | ['Row'] | ["Row(a=1, b={'world': 2, 'hello': 1})"] |
| 38 | struct_int_struct_int_str | struct<a:int,b:struct<a1:int,a2:string>> | [Row(a=1, b=Row(a1=1, a2='hello'))] | ['Row'] | ["Row(a=1, b=Row(a1=1, a2='hello'))"] |
六、按类型族的强转行为解读
6.1 整数族(tinyint / smallint / int / bigint):空值被“浮点化”
非空整数值(用例 0/2/4/6)在 legacy pandas 路径下依然以 Pythonint原样呈现,数值边界(-128/127、-32768/32767、-2147483648/2147483647、-9223372036854775808/9223372036854775807)都能无损往返。
但一旦列中出现空值(用例 1/3/5/7),行为发生显著变化:空值位置不再是None,而是被替换为float类型的nan。原因在于 legacy pandas 转换路径会把 Arrow 整数列先转换为 pandas 的Float64(可空浮点)类型,None在 pandas 中沉淀为nan,且nan属于浮点——于是 UDF 内type(x).__name__返回'float',str(x)返回'nan',而非空值也被统一转成'42.0'这样的浮点字符串。这是本 golden 文件与纯 Arrow 路径最核心的差异点(对比见第七节)。
6.2 浮点族(float / double):类型天然一致
float_values/double_values(用例 8/10)在三种路径下都是 Pythonfloat。注意测试代码使用3.14与1.0 / 3构造输入,golden 中记录的是 Spark 侧单精度/双精度存储后的值:3.140000104904175(float32 存储 3.14 的二进制近似)、0.3333333333333333(float64 存储 1/3 的近似)。float_null/double_null中空值同样变为nan,但因为列本身已是浮点,nan与正常值同属float,不会像整数列那样引入类型跳变。
6.3 Decimal:保持Decimal,空值保持None
decimal_values(用例 12)中 Spark 的Decimal('5.35')传入 UDF 后仍是Decimal类型,且精度(decimal(3,2))不丢失。decimal_null(用例 13)中空值依然是NoneType/None——这是 legacy pandas 转换路径下少数几个能完整保留None语义的类型之一(得益于Decimal对象在 pandas/Arrow 中由 Python 对象数组承载)。
6.4 字符串与二进制:无空值时直接透传
string_values(用例 14)与binary_values(用例 16)在无空值时分别以str与bytes透传,含空字符串''与空字节b'';binary_values的第三个元素在测试代码中构造为bytearray([65, 66, 67]),golden 中落为b'ABC',证明 bytearray 经 Spark 侧规整后以bytes形式到达 UDF。string_null/binary_null中的None在 pandas 2 下仍保持NoneType(这是 pandas 2 与 pandas 3 的分水岭,见第七节)。
6.5 布尔:保持bool
boolean_values/boolean_null(用例 18/19)行为稳定:非空为bool,空值为NoneType,True/False原样透传。
6.6 日期时间:Timestamp与NaTType登场
这是 legacy pandas 路径最有辨识度的一族:
date_values/date_null(用例 20/21):Sparkdate到达 UDF 后是datetime.date,空值为None,没有变化。timestamp_values(用例 22):Sparktimestamp不再以 Pythondatetime呈现,而是pandas.Timestamp(类型名'Timestamp'),字符串形式一致(2020-02-02 12:15:16.123000,微秒精度保留)。timestamp_null(用例 23):空值变成NaTType(pd.NaT,字符串形式'NaT'),非空值仍是Timestamp。
NaTType是 pandas 的“缺失时间戳”哨兵,与None判等为False且类型不同——如果在 UDF 中用if x is None过滤空时间戳,在该路径下将全部落空,必须改用pd.isna(x)或x is pd.NaT判断。
6.7 复杂类型:顶层类型稳定,嵌套整数被“numpy 化”
array/map/struct及其三层嵌套组合(用例 24–38)呈现一个统一规律:
- 顶层容器类型稳定:
array→list,map→dict,struct→Row,空容器([]、{})与空顶层值(None)行为均符合预期。 - 嵌套整数元素变成
np.int32:只要容器内部嵌有整数(如array<int>的元素、map<int,...>的键、struct中的 int 字段),元素在 legacy pandas 转换下以np.int32呈现。典型证据:array_int_null(用例 25):[4, 5, 6]显示为[np.int32(4), np.int32(5), np.int32(6)];array_array_int(用例 30):内层整数全部为np.int32(1)等;map_int_array_int(用例 33):键1/2与数组内元素均为np.int32;struct_int_array_int(用例 36):Row(a=1, b=[np.int32(1), ...]),其中a=1与b内元素均为np.int32。
- 字符串、Row、字典结构保持原样:
map<string,int>的键仍是普通str,struct内嵌的字符串字段仍是str,嵌套Row与嵌套dict不做 numpy 化。
从源码结构看,这一现象是 legacy pandas 转换路径将 Arrow 数组批量转换为pandas.Series后逐元素取出所致:整数元素被还原为 numpy 标量,而str/Row/dict等对象保持 Python 原生类型。
七、横向对比:三种路径与 pandas 主版本差异
7.1 与纯 Arrow 路径(with_arrow)的差异
对照 golden_python_udf_input_type_coercion_with_arrow.md(legacy_pandas=False),两条路径的差异可以精确归纳为“谁处理空值”和“嵌套整数给什么类型”:
| 场景 | 纯 Arrow(legacy_pandas=false) | Arrow + Legacy pandas(本文文档) |
|---|---|---|
byte_null/short_null/int_null/long_null | NoneType+None | float+nan(整数空值浮点化) |
timestamp_values | datetime | Timestamp |
timestamp_null | NoneType+datetime | NaTType(NaT)+Timestamp |
array_int_null内元素 | [4, 5, 6](普通 int) | [np.int32(4), np.int32(5), np.int32(6)] |
array_array_int内层 | [[1, 2, 3]] | [[np.int32(1), np.int32(2), np.int32(3)]] |
struct_int_array_int字段 | Row(a=1, b=[1, 2, 3]) | Row(a=1, b=[np.int32(1), ...]) |
也就是说,开启 legacy pandas 转换后:整数空值会被 pandas 提升为nan浮点、时间戳变成Timestamp/NaT、嵌套整数变成np.int32。这些差异一旦被用户代码隐式依赖(如is None判空、isinstance(x, int)类型分支),就可能造成肉眼难以发现的正确性问题——这正是该 golden 文件存在价值的最好注脚。
7.2 与 vanilla(纯 Pickle)路径的差异
对照 golden_python_udf_input_type_coercion_vanilla.md,Pickle 路径整体最“朴素”:整数族空值保持NoneType/None、时间戳是datetime、嵌套整数是普通int。也就是说 legacy pandas 路径相对默认路径引入了上述全部差异,而默认 Pickle 路径几乎不产生类型重塑。这也再次印证了 udf.py 注释中的论断——Arrow 与 Pickle 的类型强转规则不同,Arrow 优化因此默认关闭。
7.3 pandas 2 与 pandas 3 的差异:string_null用例被标记为失败
对比 pandas_3/golden_python_udf_input_type_coercion_with_arrow_and_pandas.md 可以发现,pandas 3 下仅有string_null(用例 15)一行发生改变:
| pandas 版本 | Python Type | Python Value |
|---|---|---|
| pandas 2 | ['NoneType', 'str'] | ['None', 'test'] |
| pandas 3 | ✗ Output ['nan', 'test'] != Input [None, 'test'] | (空) |
pandas 3 将字符串列中的None默认转换为nan,导致value_udf的断言assert values == input_data失败,测试将该失败本身登记为 golden 基线(✗前缀 + 空值列)。这正是测试代码pandas_dir属性注释所描述的背景:“legacy pandas 转换路径把输入经由 pandas 路由,而 pandas 3 修改了默认行为(例如字符串列的None变成nan),因此按 pandas 主版本维护独立的 golden 子目录(pandas_2/与pandas_3/),而不是在内存中打补丁”,判定逻辑见 test_python_udf_input_type.py。
对于同时需要兼容 pandas 2 与 pandas 3 的 UDF 作者,这一行数据是一个明确警示:在 legacy pandas 路径下,不要假设字符串列的空值一定是None。
八、运行测试与排查实操建议
运行本套测试(需环境具备 numpy≥2.0、pandas、pyarrow,测试类头部有
skipIf守卫):python/run-tests -k --testnames 'pyspark.sql.tests.coercion.test_python_udf_input_type'成功意味着三条路径下的 39×3 个强转结果均与 golden 基线一致。
重新生成 golden:修改了类型强转相关实现后,设置
SPARK_GENERATE_GOLDEN_FILES=1重跑上述命令即可同时刷新coercion/下的 CSV 与 Markdown(后者需安装tabulate)。UDF 内类型防御:根据本 golden 表,当用户代码开启 Arrow 优化 + legacy pandas 转换时,建议:
- 用
x is None or (isinstance(x, float) and math.isnan(x))统一处理整数列的空值(或使用pd.isna); - 判断空时间戳改用
pd.isna(x)而非x is None; - 对嵌套容器中的整数元素使用
int()归一化,避免np.int32泄漏到业务逻辑; - 兼容 pandas 3 时,对字符串列空值同时接受
None与nan。
- 用
九、总结
golden_python_udf_input_type_coercion_with_arrow_and_pandas用 39 个用例把“Arrow 优化 + Legacy pandas 转换”路径下 Python UDF 的输入强转行为完整固化了下来:标量族中整数空值被浮点化为nan、时间戳变为Timestamp/NaT,复杂类型族中嵌套整数被 numpy 化为np.int32,而字符串(pandas 2)、二进制、布尔、Decimal、date 与容器顶层类型保持稳定,pandas 3 又在字符串空值上引入了新的行为分支。这份基线既是 Spark 保证类型强转行为可回归的测试资产,也是开发者精确预判 UDF 输入类型的权威参考表——理解它,就能让 Python UDF 在不同执行路径与 pandas 版本之间写出真正健壮的代码。
- 大数据
- 数据分析
- 批处理
- 流处理
- 机器学习
- 图计算
【免费下载链接】spark
Apache Spark - A unified analytics engine for large-scale data processing
相关推荐
Apache Spark Python UDF 输入类型强制转换深度解析:Arrow 优化与 Legacy Pandas 转换路径的行为基准(pandas 3)
Apache Spark Python UDF 输入类型强制转换深度解析:Arrow 优化与 Legacy Pandas 转换路径的行为基准(pandas 3)
大数据数据分析批处理流处理机器学习图计算Apache Spark PySpark Python UDF 返回值类型强转行为全解:Arrow 优化与 Legacy Pandas 转换下的 Golden 对照矩阵
Apache Spark PySpark Python UDF 返回值类型强转行为全解:Arrow 优化与 Legacy Pandas 转换下的 Golden
大数据数据分析批处理流处理机器学习图计算深入解析 PySpark Pandas UDF 输入类型转换:golden 测试矩阵与实现原理
深入解析 PySpark Pandas UDF 输入类型转换:golden 测试矩阵与实现原理 本篇技术指南以 Apache Spark 仓库中 golden_
大数据数据分析批处理流处理机器学习图计算
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考