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

使用Boto3查询AWS Batch日志:解决动态Job ID与API报错问题

问题解决:AWS Batch日志错误检测Lambda函数优化

核心问题修复

1. describe_jobs调用失败报错

报错原因:错误使用CloudWatch Logs客户端调用AWS Batch的API。describe_jobs是Batch服务专属接口,必须用Batch客户端调用。

修复步骤:

  • 添加Batch客户端初始化:batch_client = boto3.client('batch')
  • 在get_log_stream_name函数中替换logs_client为batch_client

2. 动态获取Batch Job ID

硬编码环境变量的方式无法适配每次唯一的Job ID,提供两种可行方案:

方案A:通过Batch事件触发Lambda(推荐)

修改Lambda触发方式为Batch Job状态变化事件,此时Lambda的event参数会直接包含Job ID,无需主动查询。

CloudFormation修改点:

  • 将原定时触发事件替换为Batch事件触发器
  • 移除环境变量中的BATCH_JOB_ID配置

方案B:定时触发时主动查询最近Job

如果需要保留定时触发逻辑,调用Batch的list_jobs接口获取指定队列中最近运行的Job ID。


修正后的Lambda代码

import re
import os
import boto3
import urllib3
urllib3.disable_warnings()

# 初始化正确的客户端
logs_client = boto3.client('logs')
sns_client = boto3.client('sns')
batch_client = boto3.client('batch')  # 添加Batch客户端

def get_log_stream_name(batch_job_id):
    # 使用Batch客户端调用describe_jobs
    response = batch_client.describe_jobs(jobs=[batch_job_id])
    if response['jobs']:
        container = response['jobs'][0].get('container', {})
        log_stream_name = container.get('logStreamName')
        return log_stream_name
    else:
        print(f"No job found with ID {batch_job_id}")
        return None

def get_cloudwatch_logs(log_group_name, log_stream_name, start_time=None, end_time=None, limit=100):
    kwargs = {
        'logGroupName': log_group_name,
        'logStreamName': log_stream_name,
        'limit': limit
    }
    if start_time:
        kwargs['startTime'] = int(start_time)
    if end_time:
        kwargs['endTime'] = int(end_time)

    response = logs_client.get_log_events(**kwargs)
    return response['events']

def check_for_errors(log_events, error_keywords=None):
    if error_keywords is None:
        error_keywords = ["ERROR", "Exception", "FAILED"]
    error_pattern = re.compile("|".join(error_keywords), re.IGNORECASE)
    error_logs = [event for event in log_events if error_pattern.search(event['message'])]
    return error_logs

def send_sns_notification(topic_arn, subject, message):
    response = sns_client.publish(
        TopicArn=topic_arn,
        Subject=subject,
        Message=message
    )
    return response

def lambda_handler(event, context):
    print("event:", event)
    print("context:", context)
    
    # 方案A:从Batch事件中获取Job ID(推荐)
    batch_job_id = event.get('detail', {}).get('jobId')
    
    # 方案B:定时触发时获取最近Job(取消注释启用)
    # if not batch_job_id:
    #     response = batch_client.list_jobs(
    #         jobQueue=os.getenv('BATCH_JOB_QUEUE'),
    #         jobStatus='RUNNING'  # 可根据需求改为SUCCEEDED/FAILED
    #     )
    #     if response['jobSummaryList']:
    #         batch_job_id = sorted(response['jobSummaryList'], key=lambda x: x['createdAt'], reverse=True)[0]['jobId']
    
    log_group_name = os.getenv('LOG_GROUP_NAME')
    sns_topic_arn = os.getenv('SNS_TOPIC_ARN')  # 修正:用os.getenv替代event.getenv

    assert batch_job_id is not None, "Batch Job ID未获取到"
    assert log_group_name is not None, "日志组名称未配置"
    assert sns_topic_arn is not None, "SNS Topic ARN未配置"

    log_stream_name = get_log_stream_name(batch_job_id)
    assert log_stream_name is not None, f"无法获取Job {batch_job_id}的日志流"

    log_events = get_cloudwatch_logs(log_group_name=log_group_name, log_stream_name=log_stream_name)
    
    error_logs = check_for_errors(log_events=log_events)
    if error_logs:
        error_messages = "\n".join([f"{log['timestamp']}: {log['message']}" for log in error_logs])
        subject = f"AWS Batch任务 {batch_job_id} 日志检测到错误"
        message = f"AWS Batch任务 {batch_job_id} 日志中发现错误:\n\n{error_messages}"
        send_sns_notification(sns_topic_arn, subject, message)
        print(f"已发送SNS通知,任务 {batch_job_id} 存在错误")
    else:
        print(f"任务 {batch_job_id} 日志未检测到错误")

修正后的CloudFormation模板(方案A:Batch事件触发)

TradeFileTestLambdaFunction:
  Type: AWS::Serverless::Function
  Properties:
    Description: >
      检测AWS Batch任务日志中的错误并发送SNS通知
    Handler: app.lambda_handler
    Runtime: python3.11
    FunctionName:
      !Sub
      - '${TheAppNameForResources}-${TheEnv}'
      - TheEnv: !Ref Env
        TheAppNameForResources: !Ref AppNameForResources
    EphemeralStorage:
      Size: 10240
    Timeout: 900
    Role: !GetAtt MyLambdaExecutionRole.Arn
    Policies:
      - 'AWSLambdaVPCAccessExecutionRole'
      - 'AWSXRayDaemonWriteAccess'
      - CloudWatchLogsReadOnlyAccess
      - SNSPublishPolicy
      - Version: '2012-10-17'
        Statement:
          - Effect: Allow
            Action:
              - batch:DescribeJobs
              - batch:ListJobs  # 方案B需要保留此权限
            Resource: "*"
    Environment:
      Variables:
        ENV: !Ref Env
        SES_IDENTITY_ARN: !Ref SesIdentityArn
        SNS_TOPIC_ARN: !Ref SnsAlertTopic
        LOG_GROUP_NAME: !Ref MyBatchLogGroupName
        # 方案B需要添加队列名称
        # BATCH_JOB_QUEUE: !Ref MyBatchJobQueue
    EventInvokeConfig:
      MaximumRetryAttempts: 0
    Events:
      # 替换定时触发为Batch事件触发
      BatchJobEvent:
        Type: EventBridgeRule
        Properties:
          EventBusName: default
          Pattern:
            source:
              - "aws.batch"
            detail-type:
              - "Batch Job State Change"
            detail:
              jobQueue:
                - !Ref MyBatchJobQueue  # 指定要监控的Batch队列
              status:
                - "FAILED"  # 可根据需求添加SUCCEEDED/RUNNING等
    ReservedConcurrentExecutions: 10

SnsAlertTopic:
  Type: AWS::SNS::Topic
  Properties:
    DisplayName:
      !Sub
        - '${TheEnv}-${TheAppNameForResources}-alert-${TheRegion}'
        - TheEnv: !Ref Env
          TheAppNameForResources: !Ref AppNameForResources
          TheRegion: !Ref AWS::Region
    FifoTopic: false

# 新增Batch事件触发Lambda的权限
BatchEventPermission:
  Type: AWS::Lambda::Permission
  Properties:
    FunctionName: !Ref TradeFileTestLambdaFunction
    Action: lambda:InvokeFunction
    Principal: events.amazonaws.com
    SourceArn: !GetAtt BatchJobEvent.Arn

内容的提问来源于stack exchange,提问作者Ram

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 05:22:33