如何用Spark Scala将多查询结果追加写入单个JSON文件
解决Spark按顺序执行查询并将结果追加到同一JSON输出的问题
看起来你遇到的核心问题是:每次循环执行查询后单独调用write,导致每个查询的结果都生成独立的part文件,而不是将所有查询结果合并后统一写入(或者说整合到目标目录的输出集合中)。Spark的write.mode(Append)确实会在目标目录下新增文件,但如果每次都单独触发写入操作,就会产生多个分散的part文件——这和你想要的“同一个JSON文件”(或者说统一的输出集合)不符。
解决方案:合并所有查询结果后一次性写入
我们可以先把每个查询的结果累加合并到一个DataFrame中,最后统一执行写入操作。这样既保证了查询的顺序执行,又能将所有结果合并成指定数量的输出文件(比如通过repartition(1)生成单个文件)。
修改后的核心代码如下:
def store(jobEntity: JobDetails, jobRunId: Int): Unit = { UDFUtil.registerUdfFunctions() var outputTableName: String = null val jobQueryMap = jobEntity.jobQueryList.map(jobQuery => (jobQuery.sequenceId, jobQuery)) val sortedQueries = scala.collection.immutable.TreeMap(jobQueryMap.toSeq: _*).toMap LOGGER.debug("sortedQueries ===>" + sortedQueries) try { outputTableName = jobEntity.destinationEntity // 初始化用于合并所有结果的DataFrame var combinedDF: DataFrame = null sortedQueries.values.foreach(jobQuery => { LOGGER.debug(s"jobQuery.query ===> ${jobQuery.query}") val resultDF = SparkSession.builder.getOrCreate.sqlContext.sql(jobQuery.query) // 保留你需要的临时视图注册逻辑 tempViewsList.+=(jobQuery.queryAliasName) resultDF.createOrReplaceTempView(jobQuery.queryAliasName) // 将当前查询结果合并到总DataFrame中 combinedDF = if (combinedDF == null) { resultDF } else { // 如果所有查询的Schema一致,直接用union;如果Schema有差异,用unionByName适配 combinedDF.unionByName(resultDF, allowMissingColumns = true) } }) // 统一写入合并后的结果 if (combinedDF != null && !combinedDF.take(1).isEmpty) { val sinkDetails = new Storage(JsonUtil.toMap[String](jobEntity.sinkConnection)) val path = sinkDetails.basePath + File.separator + jobEntity.destinationEntity println("path::: " + path) // 这里repartition(1)会生成单个part文件,数据量大时建议调整分区数 combinedDF.repartition(1).write.mode(SaveMode.Append).json(path) } } catch { case e: Exception => LOGGER.error("Error storing data", e) // 保留原有的异常处理逻辑 } }
关键说明:
- 合并DataFrame:使用
unionByName而非union,可以更好地处理Schema不一致的情况(比如列顺序不同、部分列缺失);如果你的所有查询结果Schema完全一致,也可以直接用union。 - 避免多次写入:将
write操作移到循环外部,只执行一次,这样所有结果会被合并后写入目标目录,而不是每次查询都生成新的part文件。 - 分区控制:
repartition(1)会强制将所有数据合并到一个分区,从而生成单个part-00000-*.json文件。如果数据量很大,不建议这么做(会导致单个任务压力过大),可以根据实际情况调整分区数,或者去掉repartition让Spark自动处理。
如果你必须每次查询后立即追加(比如数据量极大无法全量合并):
如果因为内存或性能原因,无法将所有查询结果合并到一个DataFrame中,那可以保留循环内的写入,但需要注意:
- 去掉
repartition(1),让Spark根据集群资源自动分区,避免每次都生成单个小文件。 - 这种情况下,目标目录下还是会有多个part文件,但所有查询的结果都会追加到该目录中,符合“追加模式写入同一个JSON输出目录”的需求。
内容的提问来源于stack exchange,提问作者whatsinthename
相关产品推荐
相关产品推荐

