PySpark3.2.0下pyspark.pandas字符串编码报错问题及解决方案咨询
PySpark 3.2.0 pyspark.pandas str.encode报错问题解答
同类问题反馈情况
PySpark 3.2.x、3.3.x版本的pyspark.pandas模块由原Koalas项目合并而来,存在大量pandas原生字符串API未实现的问题,str.encode属于被广泛反馈的缺失API之一,大量从pandas迁移代码到pyspark.pandas的开发者都遇到过完全相同的NotImplementedError报错,属于版本特性未对齐的常见问题。
官方上线计划
pyspark.pandas的pandas API对齐工作随Spark大版本迭代推进,str.encode方法已经在Spark 3.4.0版本正式上线支持,3.2、3.3版本没有官方的功能回退更新计划,要使用原生的str.encode方法需要升级Spark运行环境到3.4.0及以上版本。
无需to_pandas()的替代实现方案
以下两种方案都完全在分布式集群执行,不会把全量数据拉到Driver端:
- 方案1:小数据量场景使用
apply调用原生字符串encode方法,代码改造成本最低
import pyspark.pandas as ps a = {'a':1,'b':2} df = ps.DataFrame([a]) series_columns = df.columns b = ps.Series(series_columns).apply(lambda x: x.encode('ascii', errors='ignore'))
- 方案2:大数据量场景调用Spark原生SQL的encode函数,性能更高
import pyspark.pandas as ps from pyspark.sql.functions import encode a = {'a':1,'b':2} df = ps.DataFrame([a]) series_columns = df.columns s = ps.Series(series_columns) # 转换为Spark原生DataFrame调用内置函数处理 spark_processed = s.to_spark().select(encode("value", "ascii").alias("encoded")) # 转换回pyspark.pandas Series b = spark_processed.pandas_api()["encoded"]
内容的提问来源于stack exchange,提问作者Pedro Andriow
相关产品推荐
相关产品推荐

