Spark读取含表头CSV写入Parquet因列名无效字符报错求助
解决Spark读取带表头CSV写入Parquet的列名无效字符问题
这个问题的核心原因很明确:Spark的Parquet列名不允许包含空格、 ,;{}()\n\t=这类特殊字符,而你的CSV表头里有个列名开头带空格(" M6_Debt_Review_Ind"),导致写入时触发AnalysisException错误。下面给你两种实用的解决方案,既能保留表头作为列名,又能满足Spark的列名规则:
方法一:读取后批量清理列名(通用方案)
这种方法适合你不确定所有列名是否有问题的场景,能一次性清理所有列名的无效字符:
# 读取带表头的CSV df = sql_context.read.csv("test_data_2019-01-01.csv", header=True) # 清理列名:去掉前后空格,替换其他可能的无效字符为下划线 cleaned_cols = [ col.strip() # 优先去掉前后空格,解决你遇到的开头空格问题 .replace(" ", "_") # 如果列名中间有空格,替换为下划线 .replace(",", "_") # 处理其他可能的特殊字符 .replace(";", "_") for col in df.columns ] # 重命名DataFrame的列 df_cleaned = df.toDF(*cleaned_cols) # 写入Parquet df_cleaned.write.parquet("test_data_2019-01-01.parquet")
执行这段代码后,原来的" M6_Debt_Review_Ind"会被清理成"M6_Debt_Review_Ind",其他列名也会被标准化,完全符合Spark的列名要求,写入Parquet就不会报错了。
方法二:提前定义Schema(精准控制方案)
如果已经明确知道CSV的列名和对应的数据类型,可以提前定义Schema,读取时直接使用合规的列名,同时还能提升读取性能:
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DateType # 定义Schema,列名要去掉空格、特殊字符,完全符合Spark规则 custom_schema = StructType([ StructField("foo", StringType(), nullable=True), StructField("bar", StringType(), nullable=True), StructField("bla", StringType(), nullable=True), StructField("bla2", DateType(), nullable=True), StructField("blabla", DateType(), nullable=True), StructField("bla3", IntegerType(), nullable=True), StructField("M6_Debt_Review_Ind", IntegerType(), nullable=True) ]) # 使用自定义Schema读取CSV,header=True会用来匹配表头,但列名以Schema为准 df = sql_context.read.csv( "test_data_2019-01-01.csv", header=True, schema=custom_schema ) # 直接写入Parquet df.write.parquet("test_data_2019-01-01.parquet")
这种方法的好处是不仅解决了列名问题,还能明确指定每列的数据类型,避免Spark自动推断类型可能出现的错误,适合生产环境使用。
验证结果
处理完成后,你可以通过print(df_cleaned.columns)(方法一)或print(df.columns)(方法二)查看列名,确认已经变成你期望的格式:['foo', 'bar', 'bla', 'bla2', 'blabla', 'bla3', 'M6_Debt_Review_Ind'],此时再写入Parquet就完全正常了。
内容的提问来源于stack exchange,提问作者Emile Beukes
相关产品推荐
相关产品推荐

