Synapse与Spark:导出DataFrame至CSV的文本限定符问题
问题描述
每月收到供应商提供的CSV文件,包含大量冗余列与重复数据。由于供应商需向多客户提供该文件,无法自行移除冗余列,目前正在处理重复数据问题。读取并处理文件的流程很简单:读取时用select选择所需字段(示例中省略冗余列处理),但将DataFrame导出回CSV文件时,始终无法实现为字符串列添加双引号文本限定符且不破坏输出格式的需求。
当前代码运行后,除第二列的"D,E"前后会被添加NULL导致文件加载失败外,其余功能正常。尝试移除字符串中的逗号会将"D,E"转为DE,并非可行解决方案,恳请提供解决思路。
原始数据示例
IDAppointment,IDOrganisationVisibleTo,IDPatient,DateStart 1,"A",101,2023-07-31 10:35:00.0000000 2,"B",102,2023-08-03 11:42:00.0000000 3,"C",103,2023-09-13 12:10:00.0000000 4, "D,E" ,104,2023-08-01 17:10:00.0000000
使用的代码
file_df = spark.read.options(delimiter=',', header=True, nanValue = None, inferSchema = True, quote = '"').csv(full_input_path) # 为字符串列添加双引号作为文本限定符 for col_name, dtype in file_df.dtypes: if dtype == 'string': file_df = file_df.withColumn(col_name, F.concat(F.lit('"'), F.col(col_name), F.lit('"'))) file_df = file_df.dropDuplicates(primarykey) file_df.write.option("header", "true").option("quote", "\u0000").mode("overwrite").csv(full_output_path, sep=',')
原始文件与DataFrame结构
**ORIGINAL FILE** IDAppointment,IDOrganisationVisibleTo,IDPatient,DateStart 1,"A",101,2023-07-31 10:35:00.0000000 2,"B",102,2023-08-03 11:42:00.0000000 3,"C",103,2023-09-13 12:10:00.0000000 4,"D,E",104,2023-08-01 17:10:00.0000000 root |-- IDAppointment: integer (nullable = true) |-- IDOrganisationVisibleTo: string (nullable = true) |-- IDPatient: integer (nullable = true) |-- DateStart: timestamp (nullable = true) +-------------+-----------------------+---------+-------------------+ |IDAppointment|IDOrganisationVisibleTo|IDPatient| DateStart| +-------------+-----------------------+---------+-------------------+ | 1| A| 101|2023-07-31 10:35:00| | 2| B| 102|2023-08-03 11:42:00| | 3| C| 103|2023-09-13 12:10:00| | 4| D,E| 104|2023-08-01 17:10:00| +-------------+-----------------------+---------+-------------------+
解决思路
问题根源在于手动拼接双引号后关闭了Spark的自动转义机制,导致带逗号的字符串输出异常。直接用Spark CSV Writer的原生参数即可实现需求,无需手动处理引号:
- 删除手动拼接引号的代码:不要用
concat给字符串列加双引号,Spark会自动处理需要转义的内容。 - 配置Writer的引号参数:
- 设置
quote为双引号",指定文本限定符 - 若要给所有字符串列加引号,设置
quoteAll为true;若仅需给特定列加引号,用quoteColumnList指定列名(例如"IDOrganisationVisibleTo")
- 设置
- 保留其他必要配置:header、分隔符、写入模式等参数正常设置即可。
修改后的代码示例:
file_df = spark.read.options(delimiter=',', header=True, nanValue=None, inferSchema=True, quote='"').csv(full_input_path) # 按需选择字段(示例省略,自行补充select逻辑) # file_df = file_df.select(...) file_df = file_df.dropDuplicates(primarykey) # 用Spark原生参数处理引号 file_df.write.option("header", "true") \ .option("quote", '"') \ .option("quoteAll", "true") # 替换为.option("quoteColumnList", "IDOrganisationVisibleTo")可仅给指定列加引号 .mode("overwrite") \ .csv(full_output_path, sep=',')
原理说明
之前手动加引号再关闭自动quote的方式,会让Spark输出时无法识别带逗号的字符串需要转义,导致下游解析时把"D,E"拆分为多列,进而出现NULL错误。用Spark原生的quote参数,它会自动给包含分隔符的字符串添加引号转义,同时quoteAll确保所有字符串列都被引号包裹,完全满足需求。
内容的提问来源于stack exchange,提问作者Jamie
相关产品推荐
相关产品推荐

