Azure存储文件列表获取问题求助(Scala双方案排查)
Scala实现Azure Blob/ADLS2文件元数据(路径、时间、URL)获取及问题排查
核心需求
用Scala代码实现获取Azure Blob容器(支持切换至ADLS2)内文件的路径、修改时间、访问URL。
现有方案及问题
方案1:使用SparkHadoopUtil
尝试通过SparkHadoopUtil相关方法实现,但遇到以下问题:
- 该方法属于Spark私有包,无法直接调用
- 当前POM依赖:
spark-core_2.12、spark-sql_2.123.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的providedscope即可,需提前配置Spark的Azure存储访问凭证(如通过spark.conf设置账号密钥或MSI)
针对方案2:Azure Storage Blob库卡顿问题排查与解决
本地卡顿原因
本地环境中Jackson版本冲突是核心问题:azure-storage-blob依赖特定版本Jackson,若环境存在低版本Jackson,会导致序列化/反序列化性能骤降,指定jackson-databind:2.14.2可解决版本冲突。
Databricks部署后卡顿的解决办法
检查并隔离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>优化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()切换至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
相关产品推荐
相关产品推荐

