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

Azure存储文件列表获取问题求助(Scala双方案排查)

Scala实现Azure Blob/ADLS2文件元数据(路径、时间、URL)获取及问题排查

核心需求

用Scala代码实现获取Azure Blob容器(支持切换至ADLS2)内文件的路径、修改时间、访问URL。

现有方案及问题

方案1:使用SparkHadoopUtil

尝试通过SparkHadoopUtil相关方法实现,但遇到以下问题:

  • 该方法属于Spark私有包,无法直接调用
  • 当前POM依赖:spark-core_2.12、spark-sql_2.12 3.2.2版本,scope设为provided,怀疑依赖配置是否有误

方案2:使用azure-storage-blob库

通过Azure官方的azure-storage-blob库实现,但存在卡顿问题:

  • 本地运行时,必须依赖jackson-databind:2.14.2才能避免卡顿
  • 部署至Azure Databricks后,卡顿问题依然存在,不确定是否需要对Jackson相关Jar包进行shade处理

解决方案

针对方案1:替代SparkHadoopUtil的合规方式

SparkHadoopUtil的内部方法属于私有API,不推荐直接使用,建议通过Spark公开API结合Hadoop Filesystem实现,无需依赖私有包:

import org.apache.spark.sql.SparkSession
import org.apache.hadoop.fs.{FileSystem, Path, FileStatus}

object BlobMetadataReader {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("BlobMetadataReader")
      .getOrCreate()

    val fs = FileSystem.get(spark.sparkContext.hadoopConfiguration)
    // 替换为你的Blob/ADLS2路径,格式如abfss://<container>@<account>.dfs.core.windows.net/
    val targetPath = new Path("abfss://container@storageaccount.dfs.core.windows.net/target-dir/")

    val fileStatuses: Array[FileStatus] = fs.listStatus(targetPath)
    fileStatuses.foreach { status =>
      val filePath = status.getPath.toString
      val modificationTime = status.getModificationTime
      // 根据存储类型调整URL格式,Blob用blob.core.windows.net,ADLS2用dfs.core.windows.net
      val fileUrl = s"https://storageaccount.blob.core.windows.net/container/${status.getPath.getName}"
      println(s"Path: $filePath, Modification Time: $modificationTime, URL: $fileUrl")
    }

    spark.stop()
  }
}
  • 依赖说明:保持spark-core_2.12、spark-sql_2.12的provided scope即可,需提前配置Spark的Azure存储访问凭证(如通过spark.conf设置账号密钥或MSI)

针对方案2:Azure Storage Blob库卡顿问题排查与解决

本地卡顿原因

本地环境中Jackson版本冲突是核心问题:azure-storage-blob依赖特定版本Jackson,若环境存在低版本Jackson,会导致序列化/反序列化性能骤降,指定jackson-databind:2.14.2可解决版本冲突。

Databricks部署后卡顿的解决办法

  1. 检查并隔离Jackson版本
    Databricks Runtime自带Jackson库,极易与自定义依赖冲突。先执行以下代码确认环境Jackson版本:

    println(classOf[com.fasterxml.jackson.databind.ObjectMapper].getPackage.getImplementationVersion)
    

    若版本与2.14.2不符,需通过maven-shade-plugin对Jackson包进行重命名隔离,POM配置示例:

    <plugin>
      <groupId>org.apache.maven.plugins</groupId>
      <artifactId>maven-shade-plugin</artifactId>
      <version>3.4.1</version>
      <executions>
        <execution>
          <phase>package</phase>
          <goals>
            <goal>shade</goal>
          </goals>
          <configuration>
            <relocations>
              <relocation>
                <pattern>com.fasterxml.jackson</pattern>
                <shadedPattern>your.custom.package.shaded.com.fasterxml.jackson</shadedPattern>
              </relocation>
            </relocations>
          </configuration>
        </execution>
      </executions>
    </plugin>
    
  2. 优化Blob客户端配置
    卡顿也可能是客户端连接参数不合理导致,调整连接池、超时参数提升性能:

    import com.azure.storage.blob.{BlobClientBuilder, BlobContainerClient}
    import com.azure.storage.common.StorageSharedKeyCredential
    
    val accountName = "your-storage-account"
    val accountKey = "your-storage-key"
    val containerName = "your-container"
    
    val credential = new StorageSharedKeyCredential(accountName, accountKey)
    val containerClient: BlobContainerClient = new BlobClientBuilder()
      .endpoint(s"https://$accountName.blob.core.windows.net")
      .credential(credential)
      .containerName(containerName)
      // 配置连接超时与连接池大小
      .httpClient(com.azure.core.http.netty.NettyAsyncHttpClientBuilder()
        .connectionTimeout(java.time.Duration.ofSeconds(10))
        .maxConnections(50)
        .build())
      .buildClient()
    
  3. 切换至ADLS2专用库(可选)
    若目标存储是ADLS2,建议使用azure-storage-file-datalake库,该库针对ADLS2做了性能优化:

    import com.azure.storage.file.datalake.{DataLakeFileSystemClient, DataLakeServiceClientBuilder}
    
    val serviceClient = new DataLakeServiceClientBuilder()
      .endpoint(s"https://$accountName.dfs.core.windows.net")
      .credential(credential)
      .buildClient()
    val fileSystemClient = serviceClient.getFileSystemClient(containerName)
    
    val paths = fileSystemClient.listPaths("target-dir/", true)
    paths.forEach { path =>
      val filePath = path.getName
      val modificationTime = path.getLastModified
      val fileUrl = s"https://$accountName.dfs.core.windows.net/$containerName/$filePath"
      println(s"Path: $filePath, Modification Time: $modificationTime, URL: $fileUrl")
    }
    

内容的提问来源于stack exchange,提问作者soMuchToLearnAndShare

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 17:45:40