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

VM中Apache Spark连接Azure Data Lake Gen2失败求助

虚拟机Spark连接Azure Data Lake Gen2报错排查

环境信息

  • JDK 11.0.20.1
  • Python 2.7.18
  • Spark 3.5.0
    注:此问题仅针对VM环境,与Databricks连接数据湖方法无关

问题代码

from pyspark.sql import SparkSession
from azure.identity import DefaultAzureCredential
from azure.keyvault.secrets import SecretClient

# Get Sas Token
key_vault_url = "https://<<keyvault>>.vault.azure.net/"
credential = DefaultAzureCredential() 
client = SecretClient(vault_url=key_vault_url, credential=credential)
sastoken = client.get_secret(<<SAStoken>>)

# Paths to your JAR files
path_to_hadoop_azure_jar = "/opt/spark/jars/hadoop-azure-3.3.4.jar"
path_to_azure_storage_jar = "/opt/spark/jars/azure-storage-8.6.6.jar"
path_to_jetty_util_ajax_jar = "/opt/spark/jars/jetty-util-ajax-11.0.18.jar"
path_to_jetty_util_jar = "/opt/spark/jars/jetty-util-11.0.18.jar"
path_to_azure_datalake_jar = "/opt/spark/jars/hadoop-azure-datalake-3.3.6.jar"

spark = SparkSession.builder.appName("AzureDataRead") \
    .config("spark.driver.extraClassPath", path_to_hadoop_azure_jar) \
    .config("spark.executor.extraClassPath", path_to_hadoop_azure_jar) \
    .config("spark.jars", f"{path_to_hadoop_azure_jar},{path_to_azure_storage_jar},{path_to_jetty_util_ajax_jar},{path_to_jetty_util_jar},{path_to_azure_datalake_jar}") \
    .config("fs.azure.sas.<<container>>.<<datalake>>.dfs.core.windows.net", sastoken) \
    .getOrCreate()

file_path = "/data/<<file>>" 

# Read example file
df = spark.read.format("csv") \
    .option("header", "true") \
    .option("inferSchema", "true") \
    .load(f"wasbs://<<container>>.<<datalake>>.dfs.core.windows.net/{file_path}")

# Show the DataFrame
df.show()

报错信息

Py4JJavaError: An error occurred while calling o158.load.
: java.lang.NoClassDefFoundError: Could not initialize class org.apache.hadoop.fs.azure.AzureNativeFileSystemStore
    at org.apache.hadoop.fs.azure.NativeAzureFileSystem.createDefaultStore(NativeAzureFileSystem.java:1485)
    at org.apache.hadoop.fs.azure.NativeAzureFileSystem.initialize(NativeAzureFileSystem.java:1410)
    at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:3469)
    at org.apache.hadoop.fs.FileSystem.access$300(FileSystem.java:174)
    at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:3574)
    at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:3521)
    at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:540)
    at org.apache.hadoop.fs.Path.getFileSystem(Path.java:365)
    at org.apache.spark.sql.execution.streaming.FileStreamSink$.hasMetadata(FileStreamSink.scala:53)
    at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:366)
    at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:229)
    at org.apache.spark.sql.DataFrameReader.$anonfun$load$2(DataFrameReader.scala:211)
    at scala.Option.getOrElse(Option.scala:189)
    at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:211)
    at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:186)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
    at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
    at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
    at java.base/java.lang.reflect.Method.invoke(Method.java:566)
    at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244)
    at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:374)
    at py4j.Gateway.invoke(Gateway.java:282)
    at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132)
    at py4j.commands.CallCommand.execute(CallCommand.java:79)
    at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182)
    at py4j.ClientServerConnection.run(ClientServerConnection.java:106)
    at java.base/java.lang.Thread.run(Thread.java:829)

排查与解决方法

1. 统一Hadoop相关JAR版本

当前hadoop-azure-3.3.4.jar和hadoop-azure-datalake-3.3.6.jar版本不一致,Hadoop组件必须版本统一。Spark 3.5.0适配Hadoop 3.3.4,需将hadoop-azure-datalake替换为3.3.4版本。同时azure-storage-8.6.6.jar版本过低,需替换为与Hadoop 3.3.4兼容的版本(如12.2.0)。

2. 修改文件系统协议

Azure Data Lake Gen2需使用abfss://协议,wasbs://仅适用于普通Blob Storage。修改加载路径:

load(f"abfss://<<container>>@<<datalake>>.dfs.core.windows.net/{file_path}")

3. 修正SAS Token获取方式

client.get_secret()返回的是SecretProperties对象,需提取token值:

sastoken = client.get_secret("<<SAStoken-secret-name>>").value

4. 清理冗余ClassPath配置

删除spark.driver.extraClassPath和spark.executor.extraClassPath配置,spark.jars已包含所有依赖,重复配置会导致类加载冲突。

5. 补充缺失依赖JAR

确保以下依赖存在于Spark的jars目录:

  • hadoop-common-3.3.4.jar
  • 适配Azure Storage SDK版本的azure-core-1.x.jar
  • Jackson系列JSON序列化相关JAR

修正后的示例代码

from pyspark.sql import SparkSession
from azure.identity import DefaultAzureCredential
from azure.keyvault.secrets import SecretClient

# 获取SAS Token
key_vault_url = "https://<<keyvault>>.vault.azure.net/"
credential = DefaultAzureCredential() 
client = SecretClient(vault_url=key_vault_url, credential=credential)
# 提取token值
sastoken = client.get_secret("<<SAStoken-secret-name>>").value

# 统一版本的JAR路径
path_to_hadoop_azure_jar = "/opt/spark/jars/hadoop-azure-3.3.4.jar"
path_to_azure_storage_jar = "/opt/spark/jars/azure-storage-12.2.0.jar"
path_to_jetty_util_ajax_jar = "/opt/spark/jars/jetty-util-ajax-11.0.18.jar"
path_to_jetty_util_jar = "/opt/spark/jars/jetty-util-11.0.18.jar"
path_to_azure_datalake_jar = "/opt/spark/jars/hadoop-azure-datalake-3.3.4.jar"
path_to_hadoop_common_jar = "/opt/spark/jars/hadoop-common-3.3.4.jar"

spark = SparkSession.builder.appName("AzureDataRead") \
    .config("spark.jars", f"{path_to_hadoop_azure_jar},{path_to_azure_storage_jar},{path_to_jetty_util_ajax_jar},{path_to_jetty_util_jar},{path_to_azure_datalake_jar},{path_to_hadoop_common_jar}") \
    .config("fs.azure.sas.<<container>>.<<datalake>>.dfs.core.windows.net", sastoken) \
    .getOrCreate()

file_path = "/data/<<file>>" 

# 使用abfss协议读取ADLS Gen2
df = spark.read.format("csv") \
    .option("header", "true") \
    .option("inferSchema", "true") \
    .load(f"abfss://<<container>>@<<datalake>>.dfs.core.windows.net/{file_path}")

df.show()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 01:54:51