如何在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_
相关产品推荐
相关产品推荐

