PySpark中regexp_replace未生效,求错误排查帮助
问题描述
读取的CSV内容如下:
"ZEN","123" "TEN","567"
使用PySpark的regexp_replace替换字符'E'时无效果,执行代码后没有任何替换结果,代码如下:
from pyspark.sql.functions import row_number,col,desc,date_format,to_date,to_timestamp,regexp_replace inputDirPath="/FileStore/tables/test.csv" schema = StructType() for field in fields: colType = StringType() schema.add(field.strip(),colType,True) incr_df = spark.read.format("csv").option("header", "false").schema(schema).option("delimiter", ",").option("nullValue", "").option("emptyValue","").option("multiline",True).csv(inputDirPath) for column in incr_df.columns: inc_new=incr_df.withColumn(column, regexp_replace(column,"E","") ) inc_new.show()
注:有100+列,必须用循环处理,请求排查错误。
错误原因及修复方案
核心错误:循环未累积DataFrame修改
你的循环里每次都基于原始的incr_df创建新的inc_new,而非基于上一次修改后的DataFrame。这会导致每次循环覆盖之前的修改,最终只有最后一列会被替换,其他列仍为原始值。
修复后的代码
from pyspark.sql.functions import row_number, col, desc, date_format, to_date, to_timestamp, regexp_replace from pyspark.sql.types import StructType, StringType # 补充缺失的类型导入 inputDirPath="/FileStore/tables/test.csv" # 替换为实际的列名列表,示例为["col1", "col2"] fields = ["col1", "col2"] schema = StructType() for field in fields: colType = StringType() schema.add(field.strip(), colType, True) incr_df = spark.read.format("csv")\ .option("header", "false")\ .schema(schema)\ .option("delimiter", ",")\ .option("nullValue", "")\ .option("emptyValue", "")\ .option("multiline", True)\ .csv(inputDirPath) # 初始化inc_new为原始DataFrame,累积后续修改 inc_new = incr_df for column in incr_df.columns: inc_new = inc_new.withColumn(column, regexp_replace(col(column), "E", "") ) inc_new.show()
额外注意事项
- 补充缺失导入:代码中使用了
StructType和StringType但未导入,需添加对应的类型导入语句。 - 定义
fields变量:确保fields是包含所有列名的列表,否则会触发报错。 - 显式使用
col():虽然regexp_replace可直接传列名字符串,但用col(column)能避免潜在的命名冲突,代码更清晰。
修改后,循环会依次对每一列执行替换并累积结果,所有包含'E'的列都会被正确处理。
内容的提问来源于stack exchange,提问作者Priya Chauhan
相关产品推荐
相关产品推荐

