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
相关产品推荐
相关产品推荐

