Azure Databricks读取Blob存储JSON文件遇412错误,求竞态条件解决办法
解决Azure Blob存储与Databricks管道的竞态条件问题
问题场景
使用Azure Log Analytics、Blob Storage v2、Azure Databricks Job构建数据管道:Log Analytics通过导出规则将数据写入Blob容器,Databricks挂载容器后每小时运行管道读取并转换数据。管道偶尔抛出412错误,日志如下:
Caused by: java.io.IOException: Operation failed: "The condition specified using HTTP conditional header(s) is not met.", 412, GET, https://xxx.dfs.core.windows.net/xx-xxx/WorkspaceResourceId%3D/subscriptions/xxx.json?timeout=90, ConditionNotMet, "The condition specified using HTTP conditional header(s) is not met. RequestId:xxx-xxxx-xxxx-xxxx-xxx Time:xxx-04-03T20:11:21.xxxx" at shaded.databricks.azurebfs.org.apache.hadoop.fs.azurebfs.services.AbfsInputStream.readRemote(AbfsInputStream.java:673) at shaded.databricks.azurebfs.org.apache.hadoop.fs.azurebfs.services.AbfsInputStream.readInternal(AbfsInputStream.java:619) at shaded.databricks.azurebfs.org.apache.hadoop.fs.azurebfs.services.AbfsInputStream.readOneBlock(AbfsInputStream.java:409) at shaded.databricks.azurebfs.org.apache.hadoop.fs.azurebfs.services.AbfsInputStream.read(AbfsInputStream.java:346) at java.io.DataInputStream.read(DataInputStream.java:149) at com.databricks.common.filesystem.LokiAbfsInputStream.$anonfun$read$3(LokiABFS.scala:204) at scala.runtime.java8.JFunction0$mcI$sp.apply(JFunction0$mcI$sp.java:23) at com.databricks.common.filesystem.LokiAbfsInputStream.withExceptionRewrites(LokiABFS.scala:194)
核心问题是Log Analytics导出无固定时间窗口,导致Blob存储出现读写竞态:Databricks读取时,Log Analytics仍在写入/更新目标文件,触发HTTP条件请求失败。
解决方案
1. 引入"完成标记文件"机制
- 配置Log Analytics将数据写入临时路径(如
/temp/),写入完成后手动或通过Azure Function生成对应文件名的空标记文件(如data.json.completed);若Log Analytics不支持自定义路径,可通过Blob事件网格触发Function生成标记。 - Databricks管道仅扫描带
.completed后缀的标记文件,再读取对应的源数据文件,确保读取的是已完成写入的文件。
2. 基于Blob属性判断文件稳定性
读取Blob前,检查文件的LastModifiedTime和ETag,间隔一段时间后再次校验,若两次属性一致则判定文件已停止写入:
import com.azure.storage.blob.{BlobClient, BlobContainerClientBuilder} import java.util.concurrent.TimeUnit def isBlobStable(blobClient: BlobClient, waitSeconds: Int = 30): Boolean = { val firstProps = blobClient.getProperties TimeUnit.SECONDS.sleep(waitSeconds) val secondProps = blobClient.getProperties firstProps.getLastModified.equals(secondProps.getLastModified) && firstProps.getETag.equals(secondProps.getETag) } // 使用示例 val blobClient = new BlobContainerClientBuilder() .connectionString("your-storage-connection-string") .containerName("your-container") .buildClient() .getBlobClient("path/to/target/file.json") if (isBlobStable(blobClient)) { // 执行文件读取与转换逻辑 } else { // 跳过本次处理或加入重试队列 }
3. 分层存储隔离读写
- 将Log Analytics导出目标设为临时Blob容器,通过Azure Data Factory或Blob生命周期规则,将最后修改时间超过指定阈值(如10分钟)的文件移动到只读容器。
- Databricks仅读取只读容器内的文件,彻底避免读写冲突。
4. 针对412错误添加指数退避重试
捕获412错误并实现指数退避重试,应对瞬时竞态:
import time from azure.storage.blob import BlobClient from requests.exceptions import HTTPError def read_blob_with_retry(blob_client, max_retries=3): for retry_count in range(max_retries): try: return blob_client.download_blob().readall() except HTTPError as e: if e.response.status_code == 412: time.sleep(2 ** retry_count) continue raise raise Exception(f"读取Blob失败,已重试{max_retries}次")
5. 优化Log Analytics导出分区策略
配置Log Analytics导出规则按时间(如小时)分区生成Blob文件,Databricks延迟1-2小时读取对应分区的文件,缩小读写时间窗口重叠概率。
内容的提问来源于stack exchange,提问作者inj3ct0r
相关产品推荐
相关产品推荐

