You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何按表名参数将通用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:CreateLogStream
  • logs:PutLogEvents
  • logs: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.30 14:36:03