使用Java+Spark连接Azure ADLS Gen2时进程无限停滞求助
Spark连接ADLS Gen2无限停滞的排查与修复
问题描述
尝试用Java结合Spark连接Azure ADLS Gen2存储,代码运行后仅输出少量日志便无限停滞,无任何处理动作。原始代码如下:
public static void main(String[] args) { SparkConf sparkConf = new SparkConf().setAppName("Adls connection"); SparkSession sparkSession = SparkSession.builder().config(sparkConf).master("local[*]").getOrCreate(); sparkSession.sparkContext().hadoopConfiguration().set("fs.adl.impl", "org.apache.hadoop.fs.azure.NativeAzureFileSystem"); sparkSession.sparkContext().hadoopConfiguration().set("fs.azure.account.auth.type.<account-name>.dfs.core.windows.net", "OAuth"); sparkSession.sparkContext().hadoopConfiguration().set("fs.azure.account.oauth.provider.type.<account-name>.dfs.core.windows.net", "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider"); sparkSession.sparkContext().hadoopConfiguration().set("fs.azure.account.oauth2.client.id.<account-name>.dfs.core.windows.net", "<client-id>"); sparkSession.sparkContext().hadoopConfiguration().set("fs.azure.account.oauth2.client.secret.<account-name>.dfs.core.windows.net", "<client-secret>"); sparkSession.sparkContext().hadoopConfiguration().set("fs.azure.account.oauth2.client.endpoint.<account-name>.dfs.core.windows.net", "https://login.microsoftonline.com/<tenant-id>/oauth2/token"); String path = "adl://<account-name>.azuredatalakestore.net/data/*.csv"; sparkSession.read().format("csv").option("header",true).load(path); }
排查与修复方案
1. 修正ADLS Gen2的URL协议
你误用了ADLS Gen1的adl://协议,ADLS Gen2需使用abfss://(安全传输)或abfs://协议,正确路径格式为:
abfss://<容器名称>@<存储账户名>.dfs.core.windows.net/data/*.csv
错误协议会导致Spark无法匹配对应文件系统实现,进而静默停滞。
2. 删除Gen1专属配置
代码中fs.adl.impl是ADLS Gen1的文件系统类,需直接删除该行。ADLS Gen2会自动加载对应的SecureAzureBlobFileSystem实现,无需手动配置该参数。
3. 触发Spark惰性求值
Spark的read().load()是惰性操作,仅定义读取逻辑但不会立即执行。必须调用count()、show()等动作类方法触发实际执行,否则代码停留在逻辑定义阶段,看似无限停滞。
4. 检查依赖兼容性
确保项目依赖与Spark版本匹配:
- Spark 3.x需搭配
hadoop-azure3.3.x及以上版本 - 必须包含
azure-storage-blob、azure-core等核心库
版本不匹配或依赖缺失会导致底层连接失败,且无明显报错。
5. 验证权限与认证信息
- 确认服务主体(Client ID)已被授予ADLS Gen2容器的Storage Blob Data Contributor权限(仅RBAC权限不足,需数据层面权限)
- 检查Client Secret、Tenant ID、存储账户名是否存在拼写错误,认证失败会导致Spark卡在连接阶段。
6. 开启DEBUG日志定位卡点
通过设置日志级别为DEBUG,查看停滞阶段的具体操作细节:
sparkSession.sparkContext().setLogLevel("DEBUG");
可获取认证过程、文件列表加载等关键信息,精准定位问题。
修正后的代码示例
public static void main(String[] args) { SparkConf sparkConf = new SparkConf().setAppName("AdlsGen2Connection"); SparkSession sparkSession = SparkSession.builder().config(sparkConf).master("local[*]").getOrCreate(); // 开启DEBUG日志排查问题 sparkSession.sparkContext().setLogLevel("DEBUG"); String accountName = "<account-name>"; // ADLS Gen2 OAuth配置 sparkSession.sparkContext().hadoopConfiguration().set("fs.azure.account.auth.type." + accountName + ".dfs.core.windows.net", "OAuth"); sparkSession.sparkContext().hadoopConfiguration().set("fs.azure.account.oauth.provider.type." + accountName + ".dfs.core.windows.net", "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider"); sparkSession.sparkContext().hadoopConfiguration().set("fs.azure.account.oauth2.client.id." + accountName + ".dfs.core.windows.net", "<client-id>"); sparkSession.sparkContext().hadoopConfiguration().set("fs.azure.account.oauth2.client.secret." + accountName + ".dfs.core.windows.net", "<client-secret>"); sparkSession.sparkContext().hadoopConfiguration().set("fs.azure.account.oauth2.client.endpoint." + accountName + ".dfs.core.windows.net", "https://login.microsoftonline.com/<tenant-id>/oauth2/token"); // 使用ADLS Gen2的安全协议abfss:// String path = "abfss://<container-name>@" + accountName + ".dfs.core.windows.net/data/*.csv"; // 触发实际执行 sparkSession.read().format("csv").option("header", true).load(path).show(); // 关闭SparkSession sparkSession.stop(); }
内容的提问来源于stack exchange,提问作者ashish chauhan
相关产品推荐
相关产品推荐

