ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

PySpark Python UDF 输入类型强转:Arrow 优化与 Legacy Pandas 转换路径的 Golden 基线测试深度解读

PySpark Python UDF 输入类型强转:Arrow 优化与 Legacy Pandas 转换路径的 Golden 基线测试深度解读 大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载导读本文聚焦 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 路径、纯 Picklevanilla路径的差异点以及 pandas 2 / pandas 3 主版本差异如何影响 UDF 输入的实际类型从而在编写与排查 PySpark Python UDF 时准确预判传入 Python 函数的值类型。一、背景为什么需要一份“输入类型强转”的基线测试在 PySpark 中Python UDF 的输入参数需要从 JVM 侧传入的 Spark SQL 类型转换成 Python 侧的实际对象。Spark 默认走Picklepickle-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.enabledspark.sql.legacy.execution.pythonUDF.pandas.conversion.enabledgolden 文件test_python_input_type_coercion_vanillaFalseFalsegolden_python_udf_input_type_coercion_vanillatest_python_input_type_coercion_with_arrowTrueFalsegolden_python_udf_input_type_coercion_with_arrowtest_python_input_type_coercion_with_arrow_and_pandasTrueTruepandas_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、arrayint、structa1:int,a2:string对 Python 类型使用__name__从而形成 golden 表中“Spark Type”与“Python Type”两列的统一书写。测试代码头部还给出了重新生成 golden 文件的标准命令SPARK_GENERATE_GOLDEN_FILES1 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)随后选出三路 UDFdef 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, returnTypeStringType()) value_test_udf udf(value_udf, returnTypespark_type) value_str_udf udf(value_str, returnTypeStringType()) 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 CaseSpark TypeSpark ValuePython TypePython Value0byte_valuestinyint[-128, 127, 0][int, int, int][-128, 127, 0]1byte_nulltinyint[None, 42][float, float][nan, 42.0]2short_valuessmallint[-32768, 32767, 0][int, int, int][-32768, 32767, 0]3short_nullsmallint[None, 123][float, float][nan, 123.0]4int_valuesint[-2147483648, 2147483647, 0][int, int, int][-2147483648, 2147483647, 0]5int_nullint[None, 456][float, float][nan, 456.0]6long_valuesbigint[-9223372036854775808, 9223372036854775807, 0][int, int, int][-9223372036854775808, 9223372036854775807, 0]7long_nullbigint[None, 789][float, float][nan, 789.0]8float_valuesfloat[0.0, 1.0, 3.140000104904175][float, float, float][0.0, 1.0, 3.140000104904175]9float_nullfloat[None, 3.140000104904175][float, float][nan, 3.140000104904175]10double_valuesdouble[0.0, 1.0, 0.3333333333333333][float, float, float][0.0, 1.0, 0.3333333333333333]11double_nulldouble[None, 2.71][float, float][nan, 2.71]12decimal_valuesdecimal(3,2)[Decimal(5.35), Decimal(1.23)][Decimal, Decimal][5.35, 1.23]13decimal_nulldecimal(3,2)[None, Decimal(9.99)][NoneType, Decimal][None, 9.99]14string_valuesstring[abc, , hello][str, str, str][abc, , hello]15string_nullstring[None, test][NoneType, str][None, test]16binary_valuesbinary[babc, b, bABC][bytes, bytes, bytes][babc, b, bABC]17binary_nullbinary[None, btest][NoneType, bytes][None, btest]18boolean_valuesboolean[True, False][bool, bool][True, False]19boolean_nullboolean[None, True][NoneType, bool][None, True]20date_valuesdate[datetime.date(2020, 2, 2), datetime.date(1970, 1, 1)][date, date][2020-02-02, 1970-01-01]21date_nulldate[None, datetime.date(2023, 1, 1)][NoneType, date][None, 2023-01-01]22timestamp_valuestimestamp[datetime.datetime(2020, 2, 2, 12, 15, 16, 123000)][Timestamp][2020-02-02 12:15:16.123000]23timestamp_nulltimestamp[None, datetime.datetime(2023, 1, 1, 12, 0)][NaTType, Timestamp][NaT, 2023-01-01 12:00:00]24array_int_valuesarrayint[[1, 2, 3], [], [1, None, 3]][list, list, list][[1, 2, 3], [], [1, None, 3]]25array_int_nullarrayint[None, [4, 5, 6]][NoneType, list][None, [np.int32(4), np.int32(5), np.int32(6)]]26map_str_int_valuesmapstring,int[{world: 2, hello: 1}, {}][dict, dict][{world: 2, hello: 1}, {}]27map_str_int_nullmapstring,int[None, {test: 123}][NoneType, dict][None, {test: 123}]28struct_int_str_valuesstructa1:int,a2:string[Row(a11, a2hello), Row(a12, a2world)][Row, Row][Row(a11, a2hello), Row(a12, a2world)]29struct_int_str_nullstructa1:int,a2:string[None, Row(a199, a2test)][NoneType, Row][None, Row(a199, a2test)]30array_array_intarrayarrayint[[[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)]]]31array_map_str_intarraymapstring,int[[{world: 2, hello: 1}], [{a: 1}, {b: 2}]][list, list][[{world: 2, hello: 1}], [{a: 1}, {b: 2}]]32array_struct_int_strarraystructa1:int,a2:string[[Row(a11, a2hello)], [Row(a11, a2hello), Row(a12, a2world)]][list, list][[Row(a11, a2hello)], [Row(a11, a2hello), Row(a12, a2world)]]33map_int_array_intmapint,arrayint[{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)]}]34map_int_map_str_intmapint,mapstring,int[{1: {world: 2, hello: 1}}][dict][{1: {world: 2, hello: 1}}]35map_int_struct_int_strmapint,structa1:int,a2:string[{1: Row(a11, a2hello)}][dict][{1: Row(a11, a2hello)}]36struct_int_array_intstructa:int,b:arrayint[Row(a1, b[1, 2, 3])][Row][Row(a1, b[np.int32(1), np.int32(2), np.int32(3)])]37struct_int_map_str_intstructa:int,b:mapstring,int[Row(a1, b{world: 2, hello: 1})][Row][Row(a1, b{world: 2, hello: 1})]38struct_int_struct_int_strstructa:int,b:structa1:int,a2:string[Row(a1, bRow(a11, a2hello))][Row][Row(a1, bRow(a11, a2hello))]六、按类型族的强转行为解读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__返回floatstr(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.140000104904175float32 存储 3.14 的二进制近似、0.3333333333333333float64 存储 1/3 的近似。float_null/double_null中空值同样变为nan但因为列本身已是浮点nan与正常值同属float不会像整数列那样引入类型跳变。6.3 Decimal保持Decimal空值保持Nonedecimal_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透传含空字符串与空字节bbinary_values的第三个元素在测试代码中构造为bytearray([65, 66, 67])golden 中落为bABC证明 bytearray 经 Spark 侧规整后以bytes形式到达 UDF。string_null/binary_null中的None在 pandas 2 下仍保持NoneType这是 pandas 2 与 pandas 3 的分水岭见第七节。6.5 布尔保持boolboolean_values/boolean_null用例 18/19行为稳定非空为bool空值为NoneTypeTrue/False原样透传。6.6 日期时间Timestamp与NaTType登场这是 legacy pandas 路径最有辨识度的一族date_values/date_null用例 20/21Sparkdate到达 UDF 后是datetime.date空值为None没有变化。timestamp_values用例 22Sparktimestamp不再以 Pythondatetime呈现而是pandas.Timestamp类型名Timestamp字符串形式一致2020-02-02 12:15:16.123000微秒精度保留。timestamp_null用例 23空值变成NaTTypepd.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→listmap→dictstruct→Row空容器[]、{}与空顶层值None行为均符合预期。嵌套整数元素变成np.int32只要容器内部嵌有整数如arrayint的元素、mapint,...的键、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.int32struct_int_array_int用例 36Row(a1, b[np.int32(1), ...])其中a1与b内元素均为np.int32。字符串、Row、字典结构保持原样mapstring,int的键仍是普通strstruct内嵌的字符串字段仍是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.mdlegacy_pandasFalse两条路径的差异可以精确归纳为“谁处理空值”和“嵌套整数给什么类型”场景纯 Arrowlegacy_pandasfalseArrow Legacy pandas本文文档byte_null/short_null/int_null/long_nullNoneTypeNonefloatnan整数空值浮点化timestamp_valuesdatetimeTimestamptimestamp_nullNoneTypedatetimeNaTTypeNaTTimestamparray_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(a1, b[1, 2, 3])Row(a1, 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.mdPickle 路径整体最“朴素”整数族空值保持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 TypePython Valuepandas 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_FILES1重跑上述命令即可同时刷新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 版本之间写出真正健壮的代码。赞分享大数据数据分析批处理流处理机器学习图计算【免费下载链接】sparkApache Spark - A unified analytics engine for large-scale data processing项目地址https://gitcode.com/gh_mirrors/sp/spark点击查看免费下载相关推荐Apache Spark Python UDF 输入类型强制转换深度解析Arrow 优化与 Legacy Pandas 转换路径的行为基准pandas 3Apache 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_大数据数据分析批处理流处理机器学习图计算上一篇如何快速上手SerenityOS5个简单步骤开启复古操作系统之旅下一篇告别熬夜抢茅台校园i茅台智能预约系统全攻略每天自动预约不再求人创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表