Airflow Operator内部打印日志无法在日志面板显示问题咨询
问题原因
你写的日志输出代码位于DAG定义函数内,这部分代码会在DAG文件解析阶段被调度器执行,不属于任务实例的运行时逻辑,因此不会输出到单个任务的日志页面。
除此之外你的代码还有一个直接导致配置不生效的bug:Python字典的update()方法是原地修改对象,返回值为None,你将conf.update()的返回值赋值给conf_sp,最终传给SparkSubmitOperator的conf参数实际是空值,配置自然不会生效。
不同场景的正确日志打印方式
打印DAG解析阶段的日志
这部分逻辑由DAG FileProcessor进程执行,日志不会出现在任务日志中,需要到调度器节点的DAG解析日志目录(默认路径为$AIRFLOW_HOME/logs/dagbag_manager/)或者调度器服务的标准输出中查找。
该阶段打印日志不要使用airflow.task作为logger名称,参考写法:
import logging # 直接使用当前模块的logger即可 logger = logging.getLogger(__name__) # 打印日志 logger.info("DAG解析阶段的日志内容")
打印任务运行阶段的日志
如果需要日志出现在对应任务实例的日志页面,日志代码必须写在任务的执行逻辑内:
- 对于
PythonOperator,日志代码写在python_callable指向的函数中 - 对于自定义Operator,日志代码写在重写的
execute()方法中 - 对于
SparkSubmitOperator,如果要打印提交时的配置,可以加一个前置的Python任务实现,示例代码:
import logging from datetime import datetime from airflow.operators.python import PythonOperator from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator from airflow.models import Variable from airflow import DAG def load_and_log_conf(**context): logger = logging.getLogger("airflow.task") conf = Variable.get("spark_conf", deserialize_json=True) conf.update({"spark.jars.packages":"org.apache.spark:spark-avro_2.12:2.4.3"}) logger.info(f"加载到的Spark配置为: {conf}") # 将配置推送到XCom供后续任务使用 context["ti"].xcom_push(key="spark_conf", value=conf) with DAG(dag_id="spark_demo", start_date=datetime(2024,1,1)) as dag: log_conf_task = PythonOperator( task_id="log_conf", python_callable=load_and_log_conf ) spark_task = SparkSubmitOperator( task_id="submit_spark_job", application="my_app", # 从前置任务拉取配置 conf="{{ ti.xcom_pull(task_ids='log_conf', key='spark_conf') }}" ) log_conf_task >> spark_task
这时候logger.info输出的内容就会正常显示在log_conf任务的实例日志页面中。
代码修正建议
你不需要额外定义conf_sp变量,直接修改从Variable读取的字典对象即可,修正后的配置部分代码:
conf = Variable.get("spark_conf", deserialize_json = True) conf.update({"spark.jars.packages":"org.apache.spark:spark-avro_2.12:2.4.3"}) op = SparkSubmitOperator( application = "my_app", conf = conf # 其余配置参数 )
内容的提问来源于stack exchange,提问作者user3735871
相关产品推荐
相关产品推荐

