Cloud Composer中配置结构化日志:自定义Logger替代方案求助
Google Cloud Composer 自定义结构化日志方案
由于Composer限制直接修改airflow.cfg中的logging_config_class,可以通过以下几种方式实现自定义JSON格式日志(生成jsonPayload),同时兼容GCS日志存储:
方案一:任务内双写结构化日志
在DAG任务中同时向Cloud Logging发送结构化日志,并将JSON字符串输出到Airflow默认日志,既保证Cloud Logging中有可查询的jsonPayload,又让GCS日志存储可解析的JSON文本:
import json import logging from google.cloud import logging as gcloud_logging from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime # 初始化Cloud Logging客户端 client = gcloud_logging.Client() gcloud_logger = client.logger('airflow-worker') def structured_log_task(**context): # 构造包含业务变量的结构化日志 log_content = { "message": "任务执行状态更新", "custom_key": "自定义业务值", "dag_id": context['dag'].dag_id, "task_id": context['task_instance'].task_id, "execution_date": context['execution_date'].isoformat() } # 发送到Cloud Logging,生成标准jsonPayload gcloud_logger.log_struct(log_content, severity='INFO') # 将JSON字符串输出到Airflow默认日志,GCS存储的日志将是可解析的JSON文本 logging.info(json.dumps(log_content)) with DAG( 'structured_log_demo', start_date=datetime(2024, 1, 1), schedule_interval=None, ) as dag: run_task = PythonOperator( task_id='structured_log_task', python_callable=structured_log_task, provide_context=True )
方案二:配置JSON日志格式化器(全局任务适配)
通过安装第三方日志格式化库,统一修改任务日志的输出格式,实现Cloud Logging和GCS日志的结构化:
- 在Composer环境的PyPI依赖中添加
python-json-logger - 在DAG中配置日志处理器:
import logging from pythonjsonlogger import jsonlogger from google.cloud import logging as gcloud_logging from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def setup_json_logger(): # 获取Airflow任务日志实例 task_logger = logging.getLogger('airflow.task') # 清除默认处理器,避免重复输出 for handler in task_logger.handlers[:]: task_logger.removeHandler(handler) # 配置Cloud Logging处理器,输出jsonPayload client = gcloud_logging.Client() cloud_handler = client.get_default_handler() json_formatter = jsonlogger.JsonFormatter( '%(asctime)s %(levelname)s %(dag_id)s %(task_id)s %(message)s %(extra)s' ) cloud_handler.setFormatter(json_formatter) task_logger.addHandler(cloud_handler) # 配置Stream处理器,让GCS日志存储JSON格式文本 stream_handler = logging.StreamHandler() stream_handler.setFormatter(json_formatter) task_logger.addHandler(stream_handler) return task_logger def json_log_task(**context): logger = setup_json_logger() # 通过extra传递自定义变量 logger.info("任务执行完成", extra={ "custom_key": "业务自定义字段", "execution_date": context['execution_date'].isoformat() }) with DAG( 'json_log_demo', start_date=datetime(2024, 1, 1), schedule_interval=None, ) as dag: run_task = PythonOperator( task_id='json_log_task', python_callable=json_log_task, provide_context=True )
方案三:利用Composer官方结构化日志配置(推荐)
Composer 2.x及以上版本支持通过环境变量配置结构化日志:
- 设置环境变量
AIRFLOW__LOGGING__REMOTE_LOGGING=True,启用远程日志到Cloud Logging - 设置
AIRFLOW__LOGGING__REMOTE_LOG_CONN_ID=google_cloud_default,使用默认GCP连接 - 若需要自定义格式,可将自定义
logging_config.py上传到Composer存储桶,通过环境变量AIRFLOW__LOGGING__LOGGING_CONFIG_CLASS指向该文件路径(需确认当前Composer版本是否允许此配置)
为什么之前的代码GCS日志是textPayload?
你之前的代码仅向Cloud Logging发送了结构化日志,但Airflow的默认任务日志仍以文本格式输出到GCS。上述方案通过修改任务日志输出格式或双写日志,解决了GCS日志的结构化问题。
内容的提问来源于stack exchange,提问作者Sitanshu Lamba
相关产品推荐
相关产品推荐

