Spark如何直接读取单行多JSON对象文件并生成多行DataFrame
读取单行多JSON对象文件的可行方案
原实现问题说明
你现有代码使用sc.textFile(path).collect()(0)会将整个文件内容全部拉取到Driver节点内存,仅适合测试小体积文件,生产环境遇到大文件会直接触发Driver OOM,不具备可扩展性。
方案1:分布式流式处理(全Spark版本支持,推荐)
直接在Executor端做字符串拆分和补全,全程数据不落Driver,支持任意大小的文件:
import spark.implicits._ import org.apache.spark.sql.functions._ // 读取为仅含单行文本的DataFrame val rawTextDF = spark.read.text(path) // 分布式拆分并补全每个独立JSON对象 val jsonStrDS = rawTextDF.flatMap(row => { val fullLine = row.getString(0) fullLine.split("\\}\\,\\{") .map(part => { val withHead = if (!part.startsWith("{")) s"{$part" else part if (!withHead.endsWith("}")) s"$withHead}" else withHead }) }) // 直接解析JSON字符串数据集得到结构化结果 val resultDF = spark.read.json(jsonStrDS)
方案2:利用行分隔符配置简化逻辑(Spark 3.0+支持)
通过设置自定义行分隔符直接拆分文件,再补全JSON格式即可:
import org.apache.spark.sql.functions._ val jsonStrDS = spark.read // 设置行分隔符为`},`,直接拆分出每个JSON对象主体 .option("lineSep", "},") .text(path) // 补全每个JSON对象缺失的右括号 .select( when(col("value").endsWith("}"), col("value")) .otherwise(concat(col("value"), lit("}"))) .alias("value") ) .as[String] // 解析JSON得到结构化数据 val resultDF = spark.read.json(jsonStrDS)
注意事项
如果JSON字段的value中包含},{或},的子串,上述拆分逻辑会报错,这种场景建议先将源文件预处理为标准JSON Lines格式(每个JSON对象独占一行),再用原生spark.read.json(path)读取,是兼容性最高的方案。
内容的提问来源于stack exchange,提问作者HugoDife
相关产品推荐
相关产品推荐

