如何使用PySpark批量读取多个CSV文件并快速加载到Hive分区表
解决方案
1. 批量读取目标CSV文件
Spark原生支持HDFS路径通配符匹配,可直接匹配所有FILENAMEA的CSV文件,无需逐个遍历文件:
// Scala版本,路径替换为实际HDFS路径,按需配置header、分隔符、编码等CSV参数 val csvDF = spark.read .option("header", "true") .csv("hdfs://your_hdfs_base_path/*_FILENAMEA_*.csv")
# PySpark版本 csv_df = spark.read \ .option("header", "true") \ .csv("hdfs://your_hdfs_base_path/*_FILENAMEA_*.csv")
该方式会一次性加载所有符合命名规则的文件到同一个DataFrame,比逐个读取效率提升明显。
2. 提取文件名中的ID作为分区字段
通过Spark内置的input_file_name()函数获取每行对应的源文件名,再用正则提取分区字段ID:
import org.apache.spark.sql.functions.{input_file_name, regexp_extract, col} val dfWithId = csvDF .withColumn("file_path", input_file_name()) .withColumn("ID", regexp_extract(col("file_path"), """.*/(\w+)_FILENAMEA_\d+\.csv""", 1)) .drop("file_path") // 不需要保留文件路径可直接删除
from pyspark.sql.functions import input_file_name, regexp_extract, col df_with_id = csv_df \ .withColumn("file_path", input_file_name()) \ .withColumn("ID", regexp_extract(col("file_path"), r".*/(\w+)_FILENAMEA_\d+\.csv", 1)) \ .drop("file_path")
3. 执行转换后批量写入Hive动态分区
先开启动态分区配置,无需逐个处理分区,直接按ID字段自动分区写入ORC格式Hive表:
前置配置
spark.conf.set("hive.exec.dynamic.partition", "true") spark.conf.set("hive.exec.dynamic.partition.mode", "nonstrict")
写入操作
// 先执行添加默认值等自定义转换逻辑,示例为新增默认值字段 val finalDF = dfWithId.withColumn("default_col", lit("your_default_value")) // 直接写入Hive表 finalDF.write .format("orc") .mode("append") // 全量覆盖可改为overwrite .partitionBy("ID") .saveAsTable("your_database_name.FILENAMEA") // 若必须使用HiveWarehouse Connector写入,调整格式参数即可 finalDF.write .format("com.hortonworks.spark.sql.hive.llap.HiveWarehouseConnector") .mode("append") .option("table", "your_database_name.FILENAMEA") .partitionBy("ID") .save()
4. 70张表批量处理方案
将上述逻辑封装为通用函数,传入表名作为入参,循环遍历70张表的名称列表批量执行即可,无需重复开发单表逻辑。
性能优化建议
- 调整Spark作业并行度为总文件块数的2~3倍,避免小任务过多导致调度开销过大
- 写入前对DataFrame按
ID字段重分区,减少ORC文件碎片,提升存储和后续查询效率 - 若存在大量小于HDFS块大小的CSV小文件,读取后可先执行
coalesce或repartition合并小分区,降低处理 overhead
内容的提问来源于stack exchange,提问作者Cdr
相关产品推荐
相关产品推荐

