Spark设置Azure Blob为CheckpointDir时的凭证异常及矛盾行为求助
问题:Spark设置Azure Blob检查点目录直接失败,先与存储交互后却能成功
场景说明
在Azure Synapse Spark环境中,使用org.apache.hadoop.fs.azure.AzureNativeFileSystemStore对接Azure Blob存储时,直接设置Spark检查点目录会提示缺少存储密钥,但先通过spark.readAPI和存储账户交互后,再设置检查点目录就能成功。
失败的代码示例
spark.conf.set( "fs.azure.account.key.myaccount.blob.core.windows.net", "mykey" ) spark.sparkContext.setCheckpointDir("wasbs://datalake@myaccount.blob.core.windows.net/raw/temp/")
报错信息
No credentials found for account myaccount.blob.core.windows.net in the configuration
完整PySpark堆栈信息
Py4JJavaError: An error occurred while calling o3578.setCheckpointDir. : org.apache.hadoop.fs.azure.AzureException: org.apache.hadoop.fs.azure.AzureException: No credentials found for account xxx.blob.core.windows.net in the configuration, and its container datalake is not accessible using anonymous credentials. Please check if the container exists first. If it is not publicly available, you have to provide account credentials. at org.apache.hadoop.fs.azure.AzureNativeFileSystemStore.createAzureStorageSession(AzureNativeFileSystemStore.java:1123) at org.apache.hadoop.fs.azure.AzureNativeFileSystemStore.initialize(AzureNativeFileSystemStore.java:566) at org.apache.hadoop.fs.azure.NativeAzureFileSystem.initialize(NativeAzureFileSystem.java:1423) at org.apache.hadoop.fs.FileSystem.createFileSystem(FileSystem.java:3316) at org.apache.hadoop.fs.FileSystem.access$200(FileSystem.java:137) at org.apache.hadoop.fs.FileSystem$Cache.getInternal(FileSystem.java:3365) at org.apache.hadoop.fs.FileSystem$Cache.get(FileSystem.java:3333) at org.apache.hadoop.fs.FileSystem.get(FileSystem.java:492) at org.apache.hadoop.fs.Path.getFileSystem(Path.java:361) at org.apache.spark.SparkContext.$anonfun$setCheckpointDir$2(SparkContext.scala:2595) at scala.Option.map(Option.scala:230) at org.apache.spark.SparkContext.setCheckpointDir(SparkContext.scala:2593) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) 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.GatewayConnection.run(GatewayConnection.java:238) at java.lang.Thread.run(Thread.java:750) Caused by: org.apache.hadoop.fs.azure.AzureException: No credentials found for account xxx.blob.core.windows.net in the configuration, and its container datalake is not accessible using anonymous credentials. Please check if the container exists first. If it is not publicly available, you have to provide account credentials. at org.apache.hadoop.fs.azure.AzureNativeFileSystemStore.connectUsingAnonymousCredentials(AzureNativeFileSystemStore.java:899) at org.apache.hadoop.fs.azure.AzureNativeFileSystemStore.createAzureStorageSession(AzureNativeFileSystemStore.java:1118) ... 22 more
特殊成功场景
先执行一次与目标存储账户的交互操作(即使读取的文件不存在),再设置检查点目录即可成功:
spark.read.parquet("wasbs://datalake@myaccount.blob.core.windows.net/a/b/c/invalid.parquet") spark.sparkContext.setCheckpointDir("wasbs://datalake@myaccount.blob.core.windows.net/raw/temp/")
疑问
- 这种矛盾的行为是Azure Blob存储驱动的问题,还是Spark本身的问题?
- 有没有更简便的在HDFS上设置CheckpointDir的方法?
环境相关JAR版本
/usr/hdp/current/hadoop-client/azure-storage-7.0.1.jar /usr/hdp/current/hadoop-client/hadoop-azure-3.1.1.5.0-97710309.jar
内容的提问来源于stack exchange,提问作者David Beavon
相关产品推荐
相关产品推荐

