Databricks Auto Loader写入Delta表后无法立即查询数据的问题解决
问题:Auto Loader写入Delta表后立即查询计数不更新的解决办法
我有一个简单的Auto Loader任务,代码如下:
df_dwu_limit = spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "JSON") \ .schema(schemaFromJson) \ .load("abfss://synapse-usage@xxxxx.dfs.core.windows.net/synapse-usage/")\ .writeStream \ .format("delta")\ .option("checkpointLocation", "abfss://adb-delta-synapse-usage-api@xxxxxx.dfs.core.windows.net/checkpoint_synapse_usage_api_landing/") \ .trigger(availableNow=True)\ .toTable("platform_dnu.synapse_usage_api_landing")
紧接着我就使用该Delta表进行后续处理。我在Auto Loader任务前后都执行了count(*)查询,但计数没有变化,尽管从Auto Loader的运行日志中可以看到记录已写入。如果等待约1分钟后再执行count(*)查询,就能看到更新后的记录,请问该如何解决这个问题?
以下是某一会话的Auto Loader输出:
{ "id" : "cb9a28b4-c5b4-4865-bc65-b3ca5efd2537", "runId" : "64c2afd9-ad69-4e9a-97bf-d6fa2794931a", "name" : null, "timestamp" : "2022-12-03T04:44:17.591Z", "batchId" : 7, "numInputRows" : 27, "inputRowsPerSecond" : 0.0, "processedRowsPerSecond" : 0.7879760688749453, "durationMs" : { "addBatch" : 3005, "commitOffsets" : 146, "getBatch" : 12, "latestOffset" : 30380, "queryPlanning" : 61, "triggerExecution" : 34259, "walCommit" : 222 }, "stateOperators" : [ ], "sources" : [ { "description" : "CloudFilesSource[abfss://synapse-usage@xxxxx.dfs.core.windows.net/synapse-usage/]", "startOffset" : { "seqNum" : 2534, "sourceVersion" : 1, "lastBackfillStartTimeMs" : 1669823987701, "lastBackfillFinishTimeMs" : 1669823991340 }, "endOffset" : { "seqNum" : 2562, "sourceVersion" : 1, "lastBackfillStartTimeMs" : 1669823987701, "lastBackfillFinishTimeMs" : 1669823991340 }, "latestOffset" : null, "numInputRows" : 27, "inputRowsPerSecond" : 0.0, "processedRowsPerSecond" : 0.7879760688749453, "metrics" : { "numBytesOutstanding" : "0", "numFilesOutstanding" : "0" } } ], "sink" : { "description" : "DeltaSink[abfss://adb-delta-synapse-usage-api@xxxx.dfs.core.windows.net/delta/synapse_usage_api_landing]", "numOutputRows" : -1 } }
Delta表的DDL如下:
解决方案
1. 等待流任务完全结束再执行查询
你的Auto Loader使用了.trigger(availableNow=True)模式,代码提交后只是触发任务,不会自动等待任务完成。可以调用awaitTermination()方法阻塞当前线程,直到流任务执行完毕:
df_dwu_limit = spark.readStream.format("cloudFiles") \ .option("cloudFiles.format", "JSON") \ .schema(schemaFromJson) \ .load("abfss://synapse-usage@xxxxx.dfs.core.windows.net/synapse-usage/")\ .writeStream \ .format("delta")\ .option("checkpointLocation", "abfss://adb-delta-synapse-usage-api@xxxxxx.dfs.core.windows.net/checkpoint_synapse_usage_api_landing/") \ .trigger(availableNow=True)\ .toTable("platform_dnu.synapse_usage_api_landing") # 等待流任务执行完成 df_dwu_limit.awaitTermination()
2. 强制刷新Delta表元数据
分布式存储上的Delta表元数据刷新可能有延迟,查询前执行REFRESH TABLE命令可以强制获取最新元数据:
REFRESH TABLE platform_dnu.synapse_usage_api_landing; SELECT count(*) FROM platform_dnu.synapse_usage_api_landing;
3. 确认Delta表事务提交状态
从日志中sink.numOutputRows=-1可以看出,Sink没有返回准确的输出行数,这是异步提交的正常现象。可以通过查询事务日志确认数据是否已提交:
DESCRIBE HISTORY platform_dnu.synapse_usage_api_landing LIMIT 1;
如果最新事务状态为COMMITTED,说明数据已写入完成,此时查询就能看到更新。
4. 调整元数据缓存过期时间
通过Spark配置缩短Delta表元数据的缓存TTL,让查询更快获取最新状态:
spark.conf.set("spark.databricks.delta.metadata.cacheTTLSeconds", "30")
该配置将缓存过期时间设置为30秒,减少元数据缓存导致的查询延迟。
内容的提问来源于stack exchange,提问作者Rakesh Prasad
相关产品推荐
相关产品推荐

