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
相关产品推荐
相关产品推荐

