Spark中如何从目录的数千文件里加载指定数量及匹配规则的文件
嘿,我来帮你搞定这两个Spark文件加载的问题,这在处理大量小文件的场景下确实是高频需求!下面分两种情况给你具体的实现方案:
一、加载指定数量的文件
当你需要从数千个文件里只加载前N个时,咱们可以先在Driver端获取目标目录的所有文件路径,筛选出文件(排除目录)后取指定数量,再传给Spark进行加载。这样能精准控制加载的文件数,避免Spark一次性扫描全部文件。
示例代码(Scala)
import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.spark.sql.SparkSession val spark = SparkSession.builder().appName("LoadLimitedFiles").getOrCreate() val sc = spark.sparkContext // 替换成你的目标目录路径 val dirPath = "/your/target/directory" val fs = FileSystem.get(sc.hadoopConfiguration) // 获取目录下所有文件状态,过滤掉目录,取前100个文件(可替换为你需要的数量) val filePaths = fs.listStatus(new Path(dirPath)) .filter(!_.isDirectory) .take(100) .map(_.getPath.toString) // 加载选中的文件 val df = spark.read.textFile(filePaths: _*) df.show()
示例代码(Python)
from pyspark.sql import SparkSession from py4j.java_gateway import java_import spark = SparkSession.builder.appName("LoadLimitedFiles").getOrCreate() sc = spark.sparkContext # 导入Hadoop FileSystem相关类 java_import(sc._jvm, 'org.apache.hadoop.fs.Path') java_import(sc._jvm, 'org.apache.hadoop.fs.FileSystem') # 替换成你的目标目录路径 dir_path = "/your/target/directory" fs = sc._jvm.FileSystem.get(sc._jsc.hadoopConfiguration()) # 获取文件路径列表,过滤目录,取前100个文件(可替换为你需要的数量) file_statuses = fs.listStatus(sc._jvm.Path(dir_path)) file_paths = [status.getPath().toString() for status in file_statuses if not status.isDirectory()] selected_paths = file_paths[:100] # 加载文件 df = spark.read.textFile(*selected_paths) df.show()
二、按特定命名规则加载文件
针对你的文件名规则(比如包含Japan.BAL、SelfSourcedPrivate.SHE等特定片段),有两种灵活的实现方式:
方式1:使用通配符(简单场景)
如果你的命名规则可以用通配符匹配,直接在路径里使用*就能快速筛选,这是最简便的方式:
示例:加载所有包含Japan.BAL的文件
// Scala val df = spark.read.textFile("/your/target/directory/*Japan.BAL*") // Python df = spark.read.textFile("/your/target/directory/*Japan.BAL*")
示例:同时加载包含Japan.BAL和SelfSourcedPrivate.SHE的文件
可以用逗号分隔多个通配符路径:
// Scala val df = spark.read.textFile( "/your/target/directory/*Japan.BAL*", "/your/target/directory/*SelfSourcedPrivate.SHE*" ) // Python df = spark.read.textFile( "/your/target/directory/*Japan.BAL*", "/your/target/directory/*SelfSourcedPrivate.SHE*" )
方式2:正则表达式过滤(复杂规则)
如果你的命名规则更复杂(比如需要精确匹配某几个字段的位置),可以先获取所有文件路径,再用正则表达式过滤后加载:
示例代码(Scala)
import org.apache.hadoop.fs.{FileSystem, Path} import org.apache.spark.sql.SparkSession val spark = SparkSession.builder().appName("LoadPatternFiles").getOrCreate() val sc = spark.sparkContext // 替换成你的目标目录路径 val dirPath = "/your/target/directory" val fs = FileSystem.get(sc.hadoopConfiguration) // 定义正则表达式,匹配符合规则的文件名 val pattern = """Fundamental\.FinancialLineItem\.FinancialLineItem\.(Japan\.BAL|SelfSourcedPrivate\.SHE)\..*""".r // 过滤出符合规则的文件 val filePaths = fs.listStatus(new Path(dirPath)) .filter(!_.isDirectory) .map(_.getPath.getName) .filter(pattern.matches(_)) .map(name => s"$dirPath/$name") // 加载文件 val df = spark.read.textFile(filePaths: _*) df.show()
示例代码(Python)
import re from pyspark.sql import SparkSession from py4j.java_gateway import java_import spark = SparkSession.builder.appName("LoadPatternFiles").getOrCreate() sc = spark.sparkContext java_import(sc._jvm, 'org.apache.hadoop.fs.Path') java_import(sc._jvm, 'org.apache.hadoop.fs.FileSystem') // 替换成你的目标目录路径 dir_path = "/your/target/directory" fs = sc._jvm.FileSystem.get(sc._jsc.hadoopConfiguration()) # 定义正则表达式,匹配符合规则的文件名 pattern = re.compile(r"Fundamental\.FinancialLineItem\.FinancialLineItem\.(Japan\.BAL|SelfSourcedPrivate\.SHE)\..*") # 过滤符合规则的文件 file_statuses = fs.listStatus(sc._jvm.Path(dir_path)) file_paths = [] for status in file_statuses: if not status.isDirectory(): file_name = status.getPath().getName() if pattern.match(file_name): file_paths.append(f"{dir_path}/{file_name}") # 加载文件 df = spark.read.textFile(*file_paths) df.show()
内容的提问来源于stack exchange,提问作者Anupam
相关产品推荐
相关产品推荐

