如何在PySpark中对URL格式的DataFrame列执行解码操作?
PySpark 对URL格式列做解码的实现方案
方案1:Spark 3.4+ 内置函数实现(性能最优)
Spark 3.4版本开始内置了url_decode函数,直接基于Spark引擎原生实现,没有跨语言序列化开销,性能最好,直接调用即可:
from pyspark.sql import functions as F # 直接对campaign列做URL解码,结果存入新列 df = df.withColumn("campaign_decoded", F.url_decode("campaign"))
该函数默认兼容form-urlencoded编码规则,会自动将
+转为空格、%xx格式的转义字符转为对应Unicode字符,完全匹配你的数据格式。
方案2:低版本Spark 用Pandas UDF实现(兼容所有Spark 2.3+版本)
如果你的Spark版本低于3.4,没有内置解码函数,可以用Pandas UDF封装Python的urllib解码逻辑,分布式在集群节点并行执行,性能远高于把全量数据转成pandas单机处理:
import pandas as pd from pyspark.sql import functions as F from urllib.parse import unquote_plus # 定义向量化UDF,批量处理数据,性能远高于普通逐行UDF @F.pandas_udf("string") def decode_url_udf(col: pd.Series) -> pd.Series: return col.apply(unquote_plus) df = df.withColumn("campaign_decoded", decode_url_udf("campaign"))
解码效果
针对给出的样例数据,解码后结果如下:
+--------------------+---------------------+ |user_id |campaign_decoded | +--------------------+---------------------+ |alskd9239as23093 |MM | Consórcios % | |lfifsf093039388 |Aquisição | | |kasd877191kdsd999 |Aquisição | | +--------------------+---------------------+
注意事项
- 不要使用普通的逐行Python UDF,Pandas UDF以批量为单位做数据计算和序列化,性能是普通UDF的10~100倍,完全可以支撑亿级以上大表的处理
- 所有计算均为Spark分布式执行,不需要将数据拉取到单机转pandas处理,不会出现单机内存不足、处理耗时长的问题
内容的提问来源于stack exchange,提问作者programmer987
相关产品推荐
相关产品推荐

