如何将MySQLOperator的查询输出重定向至Airflow日志?
解决Airflow中MySQLOperator查询结果不打印到日志的问题
方法一:扩展MySQLOperator自定义打印逻辑
MySQLOperator默认只会返回结果行数,不会输出具体内容。你可以继承它重写execute方法,在执行查询后把结果打印到日志:
from airflow.providers.mysql.operators.mysql import MySqlOperator import logging logger = logging.getLogger(__name__) class MySqlOperatorWithLogging(MySqlOperator): def execute(self, context): # 执行原逻辑获取结果 result = super().execute(context) # 如果有结果,打印每一行内容 if result: logger.info("查询返回的具体结果:") for row in result: logger.info(row) return result # 替换原来的MySqlOperator使用自定义的类 show_tables_master = MySqlOperatorWithLogging( task_id="show_tables_master", mysql_conn_id="dbjobs_master_mysql", sql=r"""SHOW OPEN TABLES;""", dag=dag, database="mydb", autocommit=True, )
方法二:改用PythonOperator直接执行查询并打印
如果不想自定义Operator,用PythonOperator结合MySQL驱动直接写逻辑,更灵活:
from airflow.operators.python import PythonOperator import pymysql import logging logger = logging.getLogger(__name__) def fetch_open_tables(): # 从Airflow连接获取数据库配置 from airflow.hooks.base import BaseHook conn = BaseHook.get_connection("dbjobs_master_mysql") db_conn = pymysql.connect( host=conn.host, user=conn.login, password=conn.password, database="mydb", port=conn.port or 3306 ) try: with db_conn.cursor() as cursor: cursor.execute("SHOW OPEN TABLES;") results = cursor.fetchall() logger.info("当前打开的表信息:") for row in results: logger.info(row) finally: db_conn.close() # 创建PythonOperator任务 show_tables_master = PythonOperator( task_id="show_tables_master", python_callable=fetch_open_tables, dag=dag )
两种方法都能让查询的具体结果输出到Airflow日志里,临时需求选方法二更快捷,多任务复用场景选方法一的自定义Operator更方便。
内容的提问来源于stack exchange,提问作者sFishman
相关产品推荐
相关产品推荐

