如何按表名参数将通用AWS Glue任务日志输出到不同CloudWatch日志流
按传入表名输出Glue作业日志到专属CloudWatch流的实现方案
1. 解析运行时传入的表名参数
首先在Glue作业代码开头解析表名参数,你在启动Glue作业时需要传入格式为--table_name 你的表名的运行参数,解析代码如下:
import sys import boto3 import logging import time from awsglue.utils import getResolvedOptions from botocore.exceptions import ClientError # 注册要接收的参数名 args = getResolvedOptions(sys.argv, ['JOB_NAME', 'table_name']) target_table = args['table_name']
2. 给Glue作业角色配置CloudWatch权限
你需要给Glue作业关联的IAM角色添加以下权限,资源指定为你要写入的CloudWatch日志组ARN即可:
logs:CreateLogStreamlogs:PutLogEventslogs:DescribeLogStreams
3. 替换默认日志Handler输出到自定义日志流
Glue默认的日志流是按作业名+运行ID固定生成的,你可以替换Python root logger的默认Handler,将日志直接输出到对应表名的专属日志流:
# 替换为你提前在CloudWatch创建的日志组名 LOG_GROUP_NAME = '/aws-glue/custom-table-logs' LOG_STREAM_NAME = target_table logs_client = boto3.client('logs') # 自动创建对应表名的日志流,已存在则跳过 try: logs_client.create_log_stream( logGroupName=LOG_GROUP_NAME, logStreamName=LOG_STREAM_NAME ) except ClientError as e: if e.response['Error']['Code'] != 'ResourceAlreadyExistsException': raise # 自定义CloudWatch日志输出Handler class CloudWatchLogHandler(logging.Handler): def __init__(self, log_group, log_stream): super().__init__() self.log_group = log_group self.log_stream = log_stream self.sequence_token = None def emit(self, record): log_entry = self.format(record) timestamp = int(round(time.time() * 1000)) put_params = { 'logGroupName': self.log_group, 'logStreamName': self.log_stream, 'logEvents': [{'timestamp': timestamp, 'message': log_entry}] } if self.sequence_token: put_params['sequenceToken'] = self.sequence_token try: resp = logs_client.put_log_events(**put_params) self.sequence_token = resp['nextSequenceToken'] except ClientError as e: if e.response['Error']['Code'] == 'InvalidSequenceTokenException': # 序列令牌过期时自动获取最新令牌重试 self.sequence_token = e.response['Error']['Message'].split()[-1] logs_client.put_log_events(**put_params) else: raise # 配置logger替换默认输出 logger = logging.getLogger() # *不需要覆盖Glue默认日志输出可删除接下来两行,日志会同时写入默认流和自定义流* for handler in logger.handlers[:]: logger.removeHandler(handler) # 添加自定义CloudWatch Handler cw_handler = CloudWatchLogHandler(LOG_GROUP_NAME, LOG_STREAM_NAME) cw_handler.setFormatter(logging.Formatter('%(asctime)s - %(levelname)s - %(message)s')) logger.addHandler(cw_handler) # 按需调整日志级别 logger.setLevel(logging.INFO)
可选调整
如果你需要同一张表的不同次运行日志分开存储,不用累加在同一个流里,可以将日志流名修改为f"{target_table}-{args['JOB_RUN_ID']}"按作业运行ID拆分即可,使用作业运行ID拆分流时,需要将JOB_RUN_ID添加到getResolvedOptions的参数列表中。
内容的提问来源于stack exchange,提问作者Data girl
相关产品推荐
相关产品推荐

