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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 06:35:55