Spark同一会话下动态调整写入HDFS文件副本系数的问题咨询
核心原因说明
你之前修改全局配置不生效的根本原因是:Spark Driver端修改的配置不会同步到已经启动的Executor进程,且HDFS文件写入器的配置是在作业初始化阶段读取的,运行中修改全局配置不会生效。
方案1:写入时通过作业级配置指定(最优)
Spark 3.x的DataFrame写入API支持直接通过option传入Hadoop配置项,仅对本次写入作业生效,无需修改全局配置:
// 写临时文件,副本系数1 df.write .option("dfs.replication", "1") .parquet("hdfs://path/to/tmp/file") // 写普通持久化文件,副本系数2 df.write .option("dfs.replication", "2") .parquet("hdfs://path/to/persist/file") // 高可靠场景,副本系数3 df.write .option("dfs.replication", "3") .parquet("hdfs://path/to/important/file")
如果是RDD写入场景,通过saveAsHadoopFile传入自定义的Configuration实例即可:
import org.apache.hadoop.mapred.TextOutputFormat import org.apache.hadoop.io.{Text, NullWritable} val rdd: org.apache.spark.rdd.RDD[String] = ??? // 单独为本次写入作业生成独立配置,不影响全局 val hadoopConf = rdd.sparkContext.hadoopConfiguration.newInstance() hadoopConf.set("dfs.replication", "1") rdd.map(line => (NullWritable.get(), new Text(line))) .saveAsHadoopFile( "hdfs://path/to/tmp/rddfile", classOf[NullWritable], classOf[Text], classOf[TextOutputFormat[NullWritable, Text]], hadoopConf )
方案2:写入完成后批量修改副本系数(兜底兼容)
如果现有写入逻辑无法修改,可以在写入完成后直接调用HDFS FileSystem API批量修改目标路径的副本系数,兼容性最强:
import org.apache.hadoop.fs.{FileSystem, Path} def setPathReplication(hdfsPath: String, rf: Short): Unit = { val hadoopConf = SparkSession.getActiveSession.get.sparkContext.hadoopConfiguration val fs = FileSystem.get(hadoopConf) val path = new Path(hdfsPath) // 递归修改路径下所有文件的副本系数 fs.setReplication(path, rf) } // 调用示例 df.write.parquet("hdfs://path/to/file") setPathReplication("hdfs://path/to/file", 2)
注意事项
- 无需修改
spark.sql.legacy.setCommandRejectsSparkCoreConfs配置,也不要修改全局的spark.hadoop.dfs.replication,避免全局配置混乱 - 写入时指定配置的方案无额外开销,优先使用
- 写入后修改的方案仅适合小批量路径的场景,文件量过大会给NameNode带来额外压力
内容的提问来源于stack exchange,提问作者belce
相关产品推荐
相关产品推荐

