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

Dagster中创建Telegram成功告警钩子遇错求助

Dagster Telegram成功钩子报错问题修复

初始代码及错误

Resource代码

@resource
def send_message(message):
    class TelegramConnection:
        def telegram_resource(message):
            botid = os.environ['telegram_bot']
            chat = os.environ['telegram_chat']
            bot = telegram.Bot(token=botid)   
            return bot.sendMessage(chat_id=chat, text=message, parse_mode='HTML')
    return TelegramConnection()

Hook代码

def _default_status_message(context: HookContext, status: str) -> str:
    return "Op {op_name} on job {pipeline_name} {status}!\nRun ID: {run_id}".format(
        op_name=context.op.name,
        pipeline_name=context.pipeline_name,
        run_id=context.run_id,
        status=status,
    )

def _default_success_message(context: HookContext) -> str:
    return _default_status_message(context, status="succeeded")

def telegram_on_success(
    message_fn: Callable[[HookContext], str] = _default_success_message,
    dagit_base_url: Optional[str] = "http://localhost:3000",
):

    @success_hook(required_resource_keys={"telegram"})
    def _hook(context: HookContext):
        text = message_fn(context)
        if dagit_base_url:
            text += "\n<{base_url}/instance/runs/{run_id}|View in Dagit>".format(
                base_url=dagit_base_url, run_id=context.run_id
            )
        context.resources.telegram.send_message(text=text)  # type: ignore
    return _hook

Job代码

@job (resource_defs={"telegram": send_message}, hooks={success_hook}, op_retry_policy = default_policy)
def job_text_for_pictures():
    bq = extract_ID_with_photo()
    numbers, numbers_withoun_null = extract_id_for_items()
    df_case, df_prod = extract_data_from_sf()
    subset, query = transform_two_df(df_case, df_prod, numbers, numbers_withoun_null)
    final = final_sql_result(query, subset)
    merged_df = result_merge(final, bq)
    load_df(merged_df)

初始错误

dagster._check.CheckError: Member of set mismatches type. Expected <class 'dagster._core.definitions.hook_definition.HookDefinition'>. Got <function success_hook at 0x00000284AC2BB250> of type <class 'function'>.

修改后代码及新错误

更新后的Resource代码

@resource
def telegram_resource(message):
    class TelegramConnection:
        def send_message(message):
            botid = os.environ['telegram_bot']
            chat = os.environ['telegram_chat']
            bot = telegram.Bot(token=botid)   
            return bot.sendMessage(chat_id=chat, text=message, parse_mode='HTML')
    return TelegramConnection()

更新后的Job代码

@job (resource_defs={"telegram": telegram_resource}, hooks={telegram_on_success}, op_retry_policy = default_policy)
def job_text_for_pictures():
    bq = extract_ID_with_photo()
    numbers, numbers_withoun_null = extract_id_for_items()
    df_case, df_prod = extract_data_from_sf()
    subset, query = transform_two_df(df_case, df_prod, numbers, numbers_withoun_null)
    final = final_sql_result(query, subset)
    merged_df = result_merge(final, bq)
    load_df(merged_df)

更新后的Hook代码

def _default_status_message(context: HookContext, status: str) -> str:
    return "Op {op_name} on job {pipeline_name} {status}!\nRun ID: {run_id}".format(
        op_name=context.op.name,
        pipeline_name=context.pipeline_name,
        run_id=context.run_id,
        status=status,
    )

def _default_success_message(context: HookContext) -> str:
    return _default_status_message(context, status="succeeded")

def telegram_on_success(
    message_fn: Callable[[HookContext], str] = _default_success_message,
    dagit_base_url: Optional[str] = "http://localhost:3000",
):

    @success_hook(required_resource_keys={"telegram"})
    def _hook(context: HookContext):
        text = message_fn(context)
        if dagit_base_url:
            text += "\n<{base_url}/instance/runs/{run_id}|View in Dagit>".format(
                base_url=dagit_base_url, run_id=context.run_id
            )
        context.resources.telegram.send_message(text)  # type: ignore
    return _hook

新错误

TypeError: telegram_resource..TelegramConnection.send_message() takes 1 positional argument but 2 were given

错误栈

File "C:\Users\AlBelyaev\AppData\Local\Programs\Python\Python310\lib\site-packages\dagster_core\errors.py", line 188, in user_code_error_boundary
yield
File "C:\Users\AlBelyaev\AppData\Local\Programs\Python\Python310\lib\site-packages\dagster_core\execution\plan\execute_plan.py", line 162, in _trigger_hook
hook_execution_result = hook_def.hook_fn(hook_context, step_event_list)
File "C:\Users\AlBelyaev\AppData\Local\Programs\Python\Python310\lib\site-packages\dagster_core\definitions\decorators\hook_decorator.py", line 198, in _success_hook
fn(context)
File "C:\DE\dagster\my-dagster-project\my_dagster_project\hooks\text_for_pictures.py", line 50, in _hook
context.resources.telegram.send_message(text)  # type: ignore

问题修复方案

1. 初始错误(Hook类型不匹配)修复

job装饰器的hooks参数需要传入HookDefinition实例,而非函数本身。需调用telegram_on_success()获取Hook实例后传入。

2. 新错误(send_message参数不匹配)修复

Python类的实例方法必须以self作为第一个参数,原代码中send_message缺少该参数,导致调用时自动传入的实例对象被当成第一个参数,加上传入的text后引发参数数量错误。同时,Dagster资源函数不能直接接收message参数,需遵循资源定义规范。

完整修复代码

Resource代码

import os
import telegram
from dagster import resource

@resource
def telegram_resource(context):
    class TelegramConnection:
        def __init__(self):
            self.botid = os.environ['telegram_bot']
            self.chat = os.environ['telegram_chat']
            self.bot = telegram.Bot(token=self.botid)
        
        def send_message(self, message):
            return self.bot.sendMessage(chat_id=self.chat, text=message, parse_mode='HTML')
    
    return TelegramConnection()

Hook代码

from typing import Callable, Optional
from dagster import HookContext, success_hook

def _default_status_message(context: HookContext, status: str) -> str:
    return "Op {op_name} on job {pipeline_name} {status}!\nRun ID: {run_id}".format(
        op_name=context.op.name,
        pipeline_name=context.pipeline_name,
        run_id=context.run_id,
        status=status,
    )

def _default_success_message(context: HookContext) -> str:
    return _default_status_message(context, status="succeeded")

def telegram_on_success(
    message_fn: Callable[[HookContext], str] = _default_success_message,
    dagit_base_url: Optional[str] = "http://localhost:3000",
):
    @success_hook(required_resource_keys={"telegram"})
    def _hook(context: HookContext):
        text = message_fn(context)
        if dagit_base_url:
            text += "\n<{base_url}/instance/runs/{run_id}|View in Dagit>".format(
                base_url=dagit_base_url, run_id=context.run_id
            )
        context.resources.telegram.send_message(text)
    return _hook

Job代码

from dagster import job

# 假设以下op和重试策略已提前定义
default_policy = ...
extract_ID_with_photo = ...
extract_id_for_items = ...
extract_data_from_sf = ...
transform_two_df = ...
final_sql_result = ...
result_merge = ...
load_df = ...

@job(
    resource_defs={"telegram": telegram_resource},
    hooks={telegram_on_success()},  # 调用函数获取Hook实例
    op_retry_policy=default_policy
)
def job_text_for_pictures():
    bq = extract_ID_with_photo()
    numbers, numbers_withoun_null = extract_id_for_items()
    df_case, df_prod = extract_data_from_sf()
    subset, query = transform_two_df(df_case, df_prod, numbers, numbers_withoun_null)
    final = final_sql_result(query, subset)
    merged_df = result_merge(final, bq)
    load_df(merged_df)

修复要点

  • 资源函数遵循Dagster规范,类实例方法添加self参数
  • 在类的__init__中初始化Telegram Bot,避免重复创建实例
  • Job的hooks参数传入telegram_on_success()的返回值,确保类型正确

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 01:25:18