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

Spark2.3读取未完成写入Parquet文件报错的替代解决方案咨询

可行解决方案

方案1:推动写入侧做原子写入(最高优先级)

  • HDFS同文件系统下的rename操作是原子性的,不存在短窗口读取冲突问题:只要HBase应用将重写的文件先写入同文件系统的临时目录,单文件写入完成后立即rename到正式目录,你侧读取时永远不会读到半写文件,该方案稳定性最高,改造成本也最低。

方案2:读取前预过滤合法Parquet文件(无需协调其他方,改造成本低)

Spark 2.3虽然不支持spark.sql.files.ignoreCorruptFiles参数,但你可以在触发Spark读取前,先用HDFS API遍历目标文件,仅筛选出合法的Parquet文件后再传给Spark读取,2万量级的文件校验仅需数秒即可完成,性能损耗可以忽略。
校验逻辑非常简单:合法Parquet文件的末尾4个字节固定为魔法数PAR1,仅需读取每个文件的尾部4个字节即可判断是否是写入完成的有效文件。
示例代码如下:

import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.hadoop.hdfs.HdfsConfiguration
import scala.collection.mutable.ListBuffer

val conf = spark.sparkContext.hadoopConfiguration
val fs = FileSystem.get(conf)
val validFilePaths = new ListBuffer[String]()

// 先做分区裁剪,仅遍历你需要的近3天分区路径,进一步减少校验量
val endDate = java.time.LocalDate.parse("2021-08-23")
val partitionPaths = (0 to 2).map(d => {
  val date = endDate.minusDays(d)
  s"/project/data/files/year=${date.getYear}/month=${String.format("%02d", date.getMonthValue)}/day=${String.format("%02d", date.getDayOfMonth)}"
})

partitionPaths.foreach { partitionStr =>
  val partitionPath = new Path(partitionStr)
  if (fs.exists(partitionPath)) {
    val fileStatuses = fs.listStatus(partitionPath)
    fileStatuses.filter(_.isFile).filter(_.getPath.getName.endsWith(".parquet")).foreach { fileStatus =>
      val path = fileStatus.getPath
      var stream: org.apache.hadoop.fs.FSDataInputStream = null
      try {
        val fileLen = fileStatus.getLen
        if (fileLen >= 4) {
          stream = fs.open(path)
          val magic = new Array[Byte](4)
          stream.seek(fileLen - 4)
          stream.readFully(magic)
          if (new String(magic) == "PAR1") {
            validFilePaths.append(path.toString)
          }
        }
      } catch {
        case e: Exception => // 跳过读失败的半写文件
      } finally {
        if (stream != null) stream.close()
      }
    }
  }
}

// 仅传入合法文件路径给Spark读取
val df = spark.read
  .option("spark.sql.parquet.mergeSchema", "true")
  .parquet(validFilePaths:_*)
  // 后续原有逻辑保持不变
  .select(
    $"sourcetimestamp",
    $"url",
    $"recordtype",
    $"updatetimestamp"
  )
  .withColumn("end_date", to_date(date_format(lit("2021-08-23"), "yyyy-MM-dd")))
  .withColumn("start_date", date_sub($"end_date", 2))
  .withColumn("update_date", to_date($"updatetimestamp","yyyy-MM-dd"))
  .filter( $"update_date" >= $"start_date" and $"update_date" <= $"end_date" )

方案3:自定义Parquet输入格式实现容错

如果不想提前遍历文件,可以自定义实现ParquetInputFormat,重写记录读取逻辑,遇到读文件异常时直接跳过当前文件,配置到SparkSession中即可实现和spark.sql.files.ignoreCorruptFiles相同的效果,适合需要全量扫描所有分区的场景。

方案4:优化分区裁剪减少扫描范围

你当前的代码直接读取根目录,依赖Spark自动做分区裁剪,但Spark 2.3的嵌套分区裁剪存在偶发不生效的问题,你可以根据自己的过滤条件,手动拼接需要扫描的近3天分区路径,不需要遍历24个月的全量2万个文件,不仅能大幅提升读取性能,也能降低遇到坏文件的概率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 00:15:03