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

在Prefect.io中配置任务失败邮件通知遇RuntimeError问题求助

问题分析与解决方案

错误根源

你遇到的RuntimeError是因为直接在流上下文外执行了send_failure_email('fakeemail@gmail.com'),而该函数内部调用email_send_message.submit()尝试启动Prefect任务,此时没有流上下文,触发了任务必须在流内运行的限制。

修改步骤

1. 重构failure_notification_task.py

将send_failure_email改为Prefect任务,并适配回调参数要求:

from prefect_email import email_send_message, EmailServerCredentials
from prefect import task
from dotenv import dotenv_values

@task
def send_failure_email(email_address, task, state, result):
    # task/state/result是on_failure回调自动传入的参数,必须保留
    try:
        email_server_credentials = EmailServerCredentials.load("google-email")
    except:
        env_vars = dotenv_values()
        credentials = EmailServerCredentials(
            username=env_vars.get('GOOGLE_EMAIL_ADDRESS'),
            password=env_vars.get('GOOGLE_APP_PASS'),
        )
        credentials.save("google-email", overwrite=True)
        email_server_credentials = EmailServerCredentials.load("google-email")

    # 任务内部直接调用email_send_message,无需submit
    email_send_message(
        email_server_credentials=email_server_credentials,
        subject="Pipeline Failure Notification",
        msg=f"任务 {task.name} 执行失败: {state.message}",
        email_to=email_address,
    )

@task()
def failed_task():
    raise ValueError("测试用故意抛出的错误")

2. 修正main.py的回调绑定

使用functools.partial绑定邮箱参数,传递任务引用而非直接执行函数:

from prefect import Flow
from tasks.weather_api_task import fetch_weather_api, format_responses
from tasks.azure_blob_task import connect_to_blob, get_container, upload_blob
from tasks.failure_notification_task import failed_task, send_failure_email
from dotenv import dotenv_values
from datetime import datetime
from functools import partial

# 加载环境变量
env_vars = dotenv_values()
weather_api_key = env_vars.get('WEATHER_API_KEY')
storage_account_key = env_vars.get('AZURE_STORAGE_KEY')
storage_account_name = env_vars.get('ACCOUNT_NAME')

# 配置参数
url = "http://api.weatherstack.com/current"
locations = ["Raleigh, United States", "Halifax, Canada", "Mumbai, India"]
container_name = "jistpoc"
now_str = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
blob_name = f"weather_data_{now_str}"

# 用partial绑定邮箱,传递任务引用给on_failure
failed_task.on_failure(partial(send_failure_email, 'fakeemail@gmail.com'))

# 定义流
@Flow(name="weather_flow")
def weather_flow():
    responses = fetch_weather_api(url=url, weather_api_key=weather_api_key, locations=locations)
    blob_data = format_responses(responses=responses)
    blob_service_client = connect_to_blob(storage_account_name=storage_account_name, storage_account_key=storage_account_key)
    container_client = get_container(blob_service_client=blob_service_client, container_name=container_name)
    upload_blob(container_client=container_client, blob_data=blob_data, blob_name=blob_name)
    failed_task()

# 启动流
if __name__ == "__main__":
    weather_flow.run()

关键调整说明

  • 任务化回调函数:给send_failure_email添加@task装饰器,让它能被Prefect的回调系统正确调度。
  • 避免提前执行函数:用partial绑定邮箱参数,传递任务对象而非直接调用函数,确保回调仅在任务失败时触发。
  • 适配回调参数:回调函数必须接收task、state、result三个参数,可用于在邮件中补充失败详情。
  • 标准流启动方式:用weather_flow.run()替代weather_flow._run(),符合Prefect 1.x的规范。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 01:58:27