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
相关产品推荐
相关产品推荐

