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

如何在Autoloader管道失败时将Databricks任务标记为成功?

解决方案:忽略Autoloader测试失败,确保Databricks任务标记为成功

核心问题分析

你的场景中,尽管pytest测试逻辑完全通过,但Databricks任务仍被标记为失败,原因是Spark流(Autoloader)的FAILED终止状态被任务监控系统捕获——即使Python层捕获了异常,Spark底层的流失败事件仍会被记录,触发任务失败判定。

以下是几种可行的解决方法:


方案1:主动修正失败流的终止状态

捕获query_two的异常后,主动调用stop()方法将流的终止状态从FAILED改为STOPPED,避免Databricks任务监控到失败状态。

修改后的代码示例:

# 加载预期格式数据并运行Autoloader
load_data_with_expected_schema()
query_one = run_autoloader_pipeline()
query_one.waitTermination()

# 加载含新增列的异常数据
load_data_with_unexpected_schema()
query_two_failed = False
try:
    query_two = run_autoloader_pipeline()
    query_two.awaitTermination()
except Exception as e:
    query_two_failed = True
    # 主动停止流,修正终止状态
    if query_two.isActive:
        query_two.stop()
    # 启动适配新 schema 的Autoloader
    query_three = run_autoloader_pipeline()
    query_three.awaitTermination()

# 验证测试逻辑
assert query_two_failed, "Query two未按预期失败"
# 其他断言逻辑...

# 清理所有活跃流
for query in spark.streams.active:
    query.stop()

# 明确返回成功状态
dbutils.notebook.exit(0)

方案2:将失败流的运行隔离到线程中

把query_two的执行放在独立线程内,让流的失败异常被线程捕获,不影响主线程的执行状态;同时通过共享变量验证流是否确实失败,确保测试逻辑的正确性。

代码示例:

import threading

# 共享变量传递流的失败状态
query_two_failed = False

def run_failing_query():
    global query_two_failed
    try:
        query_two = run_autoloader_pipeline()
        query_two.awaitTermination()
    except Exception as e:
        query_two_failed = True
        # 停止流(如果仍处于活跃状态)
        if 'query_two' in locals() and query_two.isActive:
            query_two.stop()

# 加载预期格式数据并运行Autoloader
load_data_with_expected_schema()
query_one = run_autoloader_pipeline()
query_one.waitTermination()

# 加载异常格式数据
load_data_with_unexpected_schema()

# 在独立线程中运行预期失败的流
thread = threading.Thread(target=run_failing_query)
thread.start()
thread.join()

# 验证流是否按预期失败
assert query_two_failed, "Query two未按预期失败"

# 启动适配新 schema 的Autoloader
query_three = run_autoloader_pipeline()
query_three.awaitTermination()

# 其他断言逻辑...

# 清理所有活跃流
for query in spark.streams.active:
    query.stop()

# 明确返回成功状态
dbutils.notebook.exit(0)

方案3:强制以pytest结果作为任务判定依据

将所有测试逻辑包裹在pytest的运行流程中,强制Databricks任务以pytest的退出码作为最终状态,忽略Spark流的中间失败事件。

代码示例:

import pytest

def test_autoloader_restart_behavior():
    # 加载预期格式数据并运行Autoloader
    load_data_with_expected_schema()
    query_one = run_autoloader_pipeline()
    query_one.waitTermination()

    # 加载含新增列的异常数据
    load_data_with_unexpected_schema()
    query_two_failed = False
    try:
        query_two = run_autoloader_pipeline()
        query_two.awaitTermination()
    except Exception as e:
        query_two_failed = True
        if query_two.isActive:
            query_two.stop()
        query_three = run_autoloader_pipeline()
        query_three.awaitTermination()

    assert query_two_failed, "Query two未按预期失败"
    # 其他断言逻辑...

# 运行测试并获取结果
test_result = pytest.main([__file__, "-v"])

# 清理所有活跃流
for query in spark.streams.active:
    query.stop()

# 根据测试结果返回任务状态
dbutils.notebook.exit(0 if test_result == 0 else 1)

内容的提问来源于stack exchange,提问作者r_g_s_

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 13:54:51