You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将Python数据清洗函数转换为PySpark实现并解决报错?

PySpark 实现多列数据清洗(支持字符串/整数类型)

问题背景

我有一个包含400+列、14000+条记录的大型DataFrame需要清洗,原本用Python编写的清洗函数在小数据量下正常,但无法适配大数据量需求,改用PySpark时出现错误:AttributeError: 'str' object has no attribute 'loc'。

原Python代码

unwanted_characters = ['[', ',', '-', '#', '@', ' '] 
cols = df.columns.to_list()

def clean_col(item):
    column= str(item.loc[col])           
    
    for character in unwanted_characters:
        if character in column:
            character_index = column.find(character)
            column = column[:character_index]
    
    return column


for x in cols:
  df[x] = lrndf.apply(clean_col, axis=1) 

尝试的PySpark代码(报错)

clean_colUDF = udf(lambda z: clean_col(z))

df.select(col("Name"), \
    convertUDF(col("Name")).alias("Name") ) \
   .show(truncate=False)

解决方案

错误原因

PySpark的UDF接收的是单个字段的原始值(字符串或整数),而非Pandas的Series对象,因此原代码中item.loc[col]的Pandas语法在PySpark中完全不适用,导致报错。

适配PySpark的实现代码

from pyspark.sql.functions import udf, col
from pyspark.sql.types import StringType, IntegerType

# 用集合存储非法字符,查找效率更高
unwanted_characters = {'[', ',', '-', '#', '@', ' '}

def clean_value(value, target_type):
    # 将任意类型值转为字符串处理
    str_val = str(value)
    
    # 找到第一个非法字符的最小索引
    first_unwanted_pos = len(str_val)
    for char in unwanted_characters:
        pos = str_val.find(char)
        if pos != -1 and pos < first_unwanted_pos:
            first_unwanted_pos = pos
    
    # 截断到第一个非法字符之前
    cleaned_str = str_val[:first_unwanted_pos]
    
    # 根据原列类型转换回对应类型
    if target_type == IntegerType():
        # 处理截断后无法转为整数的情况,返回None或按需调整
        return int(cleaned_str) if cleaned_str.isdigit() else None
    return cleaned_str

# 动态生成对应数据类型的UDF
def get_clean_udf(data_type):
    return udf(lambda x: clean_value(x, data_type), data_type)

# 遍历所有列,批量清洗
cleaned_df = df
for column_name in df.columns:
    col_type = df.schema[column_name].dataType
    # 只处理字符串和整数类型,其他类型可按需扩展
    if col_type in (StringType(), IntegerType()):
        clean_udf = get_clean_udf(col_type)
        cleaned_df = cleaned_df.withColumn(column_name, clean_udf(col(column_name)))

# 查看清洗结果
cleaned_df.show(truncate=False)

关键说明

  1. 类型兼容处理:先将所有值转为字符串执行截断逻辑,再根据原列的数据类型转换回对应类型,同时处理了整数类型截断后无法转数字的边界情况。
  2. 批量处理多列:遍历DataFrame所有列,自动识别列类型并应用对应UDF,无需手动逐个指定列。
  3. 性能优化:用集合存储非法字符,比列表的in操作效率更高,适合大数据量场景。

内容的提问来源于stack exchange,提问作者user21035178

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 12:45:50