如何让Airflow SQLExecuteQueryOperator在日志中打印SQL执行消息?
解决Airflow中SQLExecuteQueryOperator不打印PostgreSQL RAISE NOTICE消息的问题
默认情况下,Airflow的SQLExecuteQueryOperator不会捕获PostgreSQL的NOTICE级消息,要让这些消息出现在Airflow日志中,可通过以下几种方式实现:
方法一:修改Airflow数据库连接参数
在Airflow的PostgreSQL连接(conn1)配置里,给Extra字段添加JSON格式参数:
{"options": "-c client_min_messages=notice", "client_encoding": "utf8"}
该参数会让PostgreSQL客户端将NOTICE级消息返回给驱动,Airflow就能捕获并打印到日志中。
方法二:在SQLExecuteQueryOperator中指定execution_options
直接在Operator定义里添加execution_options,传递客户端消息级别参数,仅对当前任务生效:
SQLExecuteQueryOperator( task_id='t1', conn_id='conn1', sql='t1.sql', execution_options={"options": "-c client_min_messages=notice"} )
方法三:使用PostgresHook自定义执行逻辑
如果需要更灵活的控制,可通过PostgresHook手动执行脚本并捕获NOTICE消息:
from airflow.providers.postgres.hooks.postgres import PostgresHook from airflow.operators.python import PythonOperator def execute_script_with_notices(): hook = PostgresHook(postgres_conn_id='conn1') conn = hook.get_conn() # 设置客户端编码与消息级别 conn.set_client_encoding('utf8') conn.set_session(client_min_messages='notice') def notice_handler(msg): # 将NOTICE消息输出到Airflow日志 print(f"PostgreSQL NOTICE: {msg}") # 注册NOTICE消息处理器 conn.add_notice_handler(notice_handler) with conn.cursor() as cursor: with open('t1.sql', 'r') as f: sql_script = f.read() cursor.execute(sql_script) conn.commit() PythonOperator( task_id='t1', python_callable=execute_script_with_notices )
这种方式支持自定义消息处理逻辑,比如格式化输出或额外的日志记录。
内容的提问来源于stack exchange,提问作者willshen
相关产品推荐
相关产品推荐

