如何在HDFS环境下的PySpark DataFrame中处理含换行符的CSV列
解决PySpark读取含换行符CSV时的记录拆分及字段处理问题
首先得揪出你代码里的第一个关键坑:你把show()的结果赋值给了df,但show()方法只是用来打印数据,返回的是None!这就是后面所有TypeError: 'NoneType' object has no attribute '__getitem__'报错的根源。先把这个问题修正,再一步步解决换行符的问题。
第一步:正确读取含换行符的CSV文件
你的CSV里item_group_desc字段包含换行符(比如XML格式的多行内容),默认情况下Spark的CSV读取器会把换行符当成记录分隔符,直接把一条完整的记录拆成了多行。解决这个问题的核心是读取时指定multiLine=True参数,让Spark识别字段内部的换行符,把整条记录当成一个整体读取。
正确的读取代码应该是:
# 先读取数据,不要直接拼接show()!show()不返回DataFrame df = spark.read.csv( "hdfs://cluster-04d4-m/user/veerayyakumar_g/Cleansdata_Input_Test.csv", header=True, inferSchema=True, multiLine=True, # 关键参数:识别字段内的换行符 quote='"' # 如果字段是用引号包裹的,指定引号字符,确保带换行的字段被正确识别 ) # 之后单独调用show()查看数据 df.show()
加上multiLine=True后,Spark会把被引号包裹的、内部带换行的字段当成一个完整的列值,不会再拆分记录。
第二步:移除字段内的换行符
读取到正确的DataFrame后,就可以用PySpark的内置函数来清理item_group_desc列里的换行符了。注意PySpark不能像Pandas那样直接用apply(lambda x: ...)操作列,要使用Spark SQL的专用函数:
from pyspark.sql.functions import regexp_replace # 替换所有换行符(包括\n和可能的\r\n)为空字符串 new_df = df.withColumn( "item_group_desc", regexp_replace("item_group_desc", r"[\n\r]", "") ) # 查看处理后的完整结果(truncate=False不截断长字段) new_df.show(truncate=False)
这里用regexp_replace匹配所有换行符和回车符,替换为空字符串,就能把字段内的多行内容合并成单行。
完整流程示例
把两步结合起来,完整的可运行代码如下:
from pyspark.sql.functions import regexp_replace # 1. 正确读取含换行符的CSV df = spark.read.csv( "hdfs://cluster-04d4-m/user/veerayyakumar_g/Cleansdata_Input_Test.csv", header=True, inferSchema=True, multiLine=True, quote='"' ) # 2. 清理字段内的换行符 new_df = df.withColumn( "item_group_desc", regexp_replace("item_group_desc", r"[\n\r]", "") ) # 验证处理结果 new_df.select("item", "item_group_desc").show(truncate=False)
补充小提示
- 如果你的CSV字段用了其他字符包裹,记得调整
quote参数的值; inferSchema=True在大数据量下可能拖慢读取速度,如果字段类型已知,建议手动指定schema参数来提升效率。
内容的提问来源于stack exchange,提问作者Veeru Gandhad
相关产品推荐
相关产品推荐

