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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 12:32:34