You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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)
    // 保留原有的异常处理逻辑
  }
}

关键说明:

  1. 合并DataFrame:使用unionByName而非union,可以更好地处理Schema不一致的情况(比如列顺序不同、部分列缺失);如果你的所有查询结果Schema完全一致,也可以直接用union。
  2. 避免多次写入:将write操作移到循环外部,只执行一次,这样所有结果会被合并后写入目标目录,而不是每次查询都生成新的part文件。
  3. 分区控制:repartition(1)会强制将所有数据合并到一个分区,从而生成单个part-00000-*.json文件。如果数据量很大,不建议这么做(会导致单个任务压力过大),可以根据实际情况调整分区数,或者去掉repartition让Spark自动处理。

如果你必须每次查询后立即追加(比如数据量极大无法全量合并):

如果因为内存或性能原因,无法将所有查询结果合并到一个DataFrame中,那可以保留循环内的写入,但需要注意:

  • 去掉repartition(1),让Spark根据集群资源自动分区,避免每次都生成单个小文件。
  • 这种情况下,目标目录下还是会有多个part文件,但所有查询的结果都会追加到该目录中,符合“追加模式写入同一个JSON输出目录”的需求。

内容的提问来源于stack exchange,提问作者whatsinthename

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.08 10:22:45