PySpark中如何按字节大小正确过滤DataFrame记录
PySpark按单条记录字节大小过滤DataFrame的异常修复
现有代码核心问题
- 逻辑冗余且引用错误:两次计算size列时均引用未新增
s3Key字段的updatedFrame,而非字段更新后的dfFinal,会导致列计算范围不符合预期,甚至出现列不存在的报错 - 字节数计算逻辑错误:传入UDF的参数已经是
to_json生成的JSON字符串,UDF内部又做了一次json.dumps,会额外转义引号、转义符等内容,导致计算出的字节数远大于实际值;遇到null值、特殊编码字符时还会直接抛空指针或编码异常 - 性能差:Python UDF需要在JVM和Python进程间做跨进程数据序列化,百万级数据下性能损耗远高于Spark内置函数
- 未做边界兼容:遇到二进制类型、嵌套复杂结构(多层Map/Array)时,
json.dumps会直接抛出序列化异常
可直接运行的修复方案
优先用Spark内置函数替代自定义UDF,无序列化问题、性能是Python UDF的5倍以上:
from pyspark.sql import functions as F from pyspark.sql.types import LongType # 原有TTL字段逻辑保留 purge_timestamp = getPurgeTimestamp() updatedFrame = df.withColumn('ttlExpDate', F.lit(purge_timestamp).cast("long")) # 所有业务字段(包括s3Key)全部新增完成后,再统一计算记录大小,不要重复计算 dfFinal = updatedFrame.withColumn('s3Key', F.lit("")) # 一次性计算整行记录的UTF-8编码字节数,单位KB all_columns = dfFinal.columns dfFinal = dfFinal.withColumn( "size_kb", (F.length(F.to_json(F.struct(*[F.col(c) for c in all_columns])).encode("UTF-8")) / 1024).cast(LongType()) ) # 按需求过滤,比如过滤小于等于1MB的记录 filtered_result = dfFinal.filter(F.col("size_kb") <= 1024).drop("size_kb")
特殊场景适配
- 如果需要计算数据实际存储的字节大小(而非JSON序列化后的大小),可以直接调用Spark内置的大小计算表达式,精度更高:
# 注册内置大小计算函数 spark.udf.register("get_raw_size", "org.apache.spark.sql.catalyst.expressions.UsedSize", LongType()) dfFinal = dfFinal.withColumn( "size_kb", (F.expr(f"get_raw_size(struct({','.join(all_columns)}))") / 1024).cast(LongType()) )
- 如果必须保留原有UDF逻辑,需要修复重复序列化的bug,同时增加空值兼容:
from pyspark.sql.functions import udf from pyspark.sql.types import IntegerType @udf(returnType=IntegerType()) def getSize(json_str): if not json_str: return 0 return len(json_str.encode('utf-8', 'ignore')) // 1024 # 调用时直接传入to_json结果即可,不需要在UDF内重复做JSON序列化 dfFinal = dfFinal.withColumn('size', getSize(F.to_json(F.struct(*[F.col(x) for x in all_columns]))))
大数据量排查技巧
百万级数据不需要全量打印日志,用以下方式快速定位问题:
- 用
dfFinal.sample(fraction=0.001).show(truncate=False)抽取0.1%的样本查看计算结果,数据量极小可以直接输出 - 直接查看Spark UI的SQL页面,每个计算阶段的报错、异常记录类型都会明确展示,不需要依赖本地打印日志
- 遇到序列化失败的字段,先单独对二进制、复杂嵌套类型做转base64、转字符串处理后再参与大小计算
内容的提问来源于stack exchange,提问作者Francisco Pereira
相关产品推荐
相关产品推荐

