使用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
相关产品推荐
相关产品推荐

