使用PySpark写入CSV时如何保留字符串列中的\r\n(不换行)
解决Spark写入CSV时保留字符串内\r\n字面量的问题
方案一:Spark内置函数替换(优先推荐)
直接通过Spark内置的regexp_replace函数,将字符串中的实际换行控制字符(\r、\n)替换为字面量形式的\r和\n,避免CSV写入器将其解析为记录换行。
Python代码示例
from pyspark.sql import functions as F # 替换columnA中的\r和\n为字面量形式 processed_df = df2.withColumn("columnA", F.regexp_replace(F.regexp_replace(F.col("columnA"), "\r", "\\\\r"), "\n", "\\\\n") ) # 写入CSV,保持原有配置 processed_df.coalesce(1).write.options( delimiter='|', quoteAll=True, escape='\\' ).csv("/target/output/path")
Scala代码示例
import org.apache.spark.sql.functions._ val processedDf = df2.withColumn("columnA", regexp_replace(regexp_replace(col("columnA"), "\r", "\\\\r"), "\n", "\\\\n") ) processedDf.coalesce(1).write.options( Map("delimiter" -> "|", "quoteAll" -> "true", "escape" -> "\\") ).csv("/target/output/path")
方案二:自定义UDF处理(复杂场景适配)
如果需要更灵活的字符串处理逻辑,可以编写自定义UDF:
from pyspark.sql.functions import udf from pyspark.sql.types import StringType from pyspark.sql import functions as F def escape_line_breaks(s): if s is None: return s # 将实际换行符替换为字面量 return s.replace("\r", "\\r").replace("\n", "\\n") # 注册UDF escape_udf = udf(escape_line_breaks, StringType()) processed_df = df2.withColumn("columnA", escape_udf(F.col("columnA"))) processed_df.coalesce(1).write.options( delimiter='|', quoteAll=True, escape='\\' ).csv("/target/output/path")
方案三:Pandas备选方案
若Spark方案仍不符合需求,可转用Pandas处理:
# 将Spark DataFrame转为Pandas DataFrame pandas_df = df2.toPandas() # 替换换行符为字面量 pandas_df["columnA"] = pandas_df["columnA"].str.replace("\r", "\\r").str.replace("\n", "\\n") # 写入CSV,配置匹配需求 pandas_df.to_csv( "/target/output/file.csv", sep="|", quotechar='"', quoting=1, # 对应QUOTE_ALL,所有字段加引号 escapechar="\\", index=False # 不写入索引列 )
注意事项
- 验证输出时请使用纯文本编辑器(如Notepad++)查看原始内容,Excel等表格工具可能会自动解析字面量
\r\n为实际换行。 coalesce(1)用于生成单个输出文件,可根据实际集群规模调整分区数。
内容的提问来源于stack exchange,提问作者PipelineSurfer
相关产品推荐
相关产品推荐

