Spark Streaming读取ADLS Gen2文本文件报错及解决方案咨询
解决方案
1. 修正流读取路径(核心问题)
Spark Structured Streaming 的 text 数据源不支持单个文件的流式读取,它要求路径指向一个目录,用于监控目录内新增的文件。你的代码中使用 /*.txt 匹配文件,不符合流处理的工作模式,这是导致令牌错误和无法处理单个文件的核心原因之一。
- 调整路径为目录路径:
file_path = "abfss://<container>@<account-name>.dfs.core.windows.net/<path>/" - 如果仅需处理单个文件,改用批处理而非流处理:
# 批处理读取单个文件 batch_df = spark.read.schema(schema).text("abfss://<container>@<account-name>.dfs.core.windows.net/<path>/single-file.txt") # 后续处理逻辑 batch_df_transformed.show()
2. 切换到ADLS Gen2推荐的ABFSS协议
避免使用旧的 adl:// 协议,改用ADLS Gen2官方推荐的 abfss:// 协议,减少协议兼容带来的权限或令牌问题。
3. 正确配置SAS令牌权限
如果使用SAS令牌,确保全局配置Spark参数,让流处理的所有环节(包括检查点目录)都能获取到令牌:
# 设置SAS令牌全局配置 spark.conf.set("fs.azure.account.auth.type.<account-name>.dfs.core.windows.net", "SAS") spark.conf.set("fs.azure.sas.token.provider.type.<account-name>.dfs.core.windows.net", "org.apache.hadoop.fs.azurebfs.sas.FixedSASTokenProvider") spark.conf.set("fs.azure.sas.fixed.token.<account-name>.dfs.core.windows.net", "<your-full-sas-token>") # 流读取目录 streaming_df = spark.readStream \ .schema(schema) \ .text("abfss://<container>@<account-name>.dfs.core.windows.net/<path>/") # 定义流查询时必须指定检查点目录 query = streaming_df_transformed.writeStream \ .outputMode("append") \ .format("console") \ .option("checkpointLocation", "abfss://<container>@<account-name>.dfs.core.windows.net/<checkpoint-path>/") \ .start() query.awaitTermination()
4. 确保检查点目录权限
流处理必须指定 checkpointLocation,该目录需要ADLS Gen2的读写权限。如果未指定,Spark会尝试使用默认路径,可能因权限不足触发令牌错误。
内容的提问来源于stack exchange,提问作者Sasidhar Reddy
相关产品推荐
相关产品推荐

