Spark Scala读写Kafka流写入CSV时特殊字符清理及报错排查
Spark Scala 处理Kafka到CSV写入字段错位问题解决方案
问题现象
基于Spark Scala实现Kafka流数据读取写入CSV的场景下,因字段内包含换行符、制表符、特殊Unicode字符,导致输出CSV内容被拆分为多行、字段错位,后续建表异常。
原始JSON格式数据样例:
{"Request ID": "XXXX-XXXX", "Owner": "cam,Ash (acf)", "Request Description": "nasvote.", "Short Term Resolution": "‑322\tCRT cod. \n∣026822\tMa (completed 8/18)\no\tACDC confirmed materials pacted\nⷰ02322\tClaims Re1.\no\t"}
异常输出的CSV样例:
XXXX-XXXX,"cam,Ash (acf)",nasvote.,• CRT cod. • Ma (completed 8/18) o ACDC confirmed materials pacted • Claims Re1. o
调试过程中两类典型报错:
- 调用
regexp_replace传入列参数时触发ArrayOutOfBound异常 - 自定义Schema解析JSON时抛出字段不存在异常:
org.apache.spark.sql.AnalysisException: cannot resolve '`Short Term Resolution`' given input columns: [key, jsontostructs(CAST(value AS STRING))];
分步解决方法
1. 正确解析JSON,解决字段找不到报错
上述Schema解析报错的核心原因是:from_json解析后生成的是结构体类型列,未展开结构体就直接引用内部字段,自然无法匹配到列。正确处理逻辑如下:
首先定义和JSON字段完全匹配的Schema:
import org.apache.spark.sql.types._ // 字段名必须和JSON中的key完全一致,包括空格 val jsonSchema = StructType(Seq( StructField("Request ID", StringType, nullable = true), StructField("Owner", StringType, nullable = true), StructField("Request Description", StringType, nullable = true), StructField("Short Term Resolution", StringType, nullable = true) ))
解析Kafka二进制数据并展开结构体:
import org.apache.spark.sql.functions._ val parsedDf = kafkaRawDf // Kafka读取的value默认是二进制类型,先转字符串 .withColumn("value_str", col("value").cast("string")) // 按定义的Schema解析JSON为结构体列 .withColumn("json_data", from_json(col("value_str"), jsonSchema)) // 展开结构体所有字段,得到平级的业务列 .select("key", "json_data.*")
执行完这一步即可直接引用Short Term Resolution这类带空格的字段,不会再报字段不存在的错误。
2. 批量清理特殊字符,解决行拆分问题
regexp_replace触发数组越界,大多是因为正则规则写法错误、或未按Column类型传参导致。不要写复杂嵌套正则,按字符类型分层替换即可,逻辑稳定不易出错:
val cleanedDf = parsedDf // 对所有字符串类型字段统一做清理 .select(parsedDf.columns.map { colName => val processCol = regexp_replace( regexp_replace( regexp_replace(col(colName), "[\\n\\r]", " "), // 换行符替换为空格 "\\t", " " // 制表符替换为空格 ), // 移除ASCII可见范围外的不可见控制字符,保留常规英文、数字、标点 "[\\p{Cntrl}&&[^\u0020-\u007E]]", "" ) processCol.as(colName) }: _*)
如果业务需要保留中文、特殊标点等非ASCII字符,可以调整最后一段正则的匹配范围,仅移除会导致CSV换行的控制类字符即可。
3. 配置CSV写入参数,兜底避免错位
字符清理不能覆盖100%异常场景,必须在写入层配置转义规则做双重保障:
cleanedDf.writeStream .format("csv") .option("path", "/your/output/path") .option("header", "true") .option("quoteAll", "true") // 所有字段值用双引号包裹 .option("escape", "\"") // 字段内的双引号做转义处理 .option("multiLine", "false") // 禁止写出多行内容 .start()
避坑说明
- 引用带空格的字段名时需要用反引号包裹,结构体展开后可直接引用,不需要额外转义
regexp_replace第一个参数必须是col()包装的Column类型对象,直接传字符串列名会触发类型匹配错误- 不要仅依赖字符清理解决CSV错位问题,必须配合写入端的引号、转义配置做兜底
内容的提问来源于stack exchange,提问作者Ashwanth arul
相关产品推荐
相关产品推荐

