SQS触发Lambda读取S3对象遇TypeError报错求助
解决SQS触发Lambda读取S3日志并导入OpenSearch的问题
问题根源
你碰到的TypeError: string indices must be integers,核心原因是从SQS消息中提取的record['body']是JSON格式字符串,而非Python字典,直接用字典索引访问['Records']会报错。你的代码仅做了赋值操作,没有完成字符串到JSON对象的解析。
解决方案步骤
- 导入
json模块,将SQS消息体的JSON字符串解析为Python字典 - 处理S3对象键的URL编码(比如空格、特殊字符会被SQS转义,需解码还原)
- 读取S3中的Apache日志文件,逐行解析日志格式
- 将解析后的日志批量写入OpenSearch索引
完整修正代码
import json import boto3 import urllib.parse from opensearchpy import OpenSearch, RequestsHttpConnection from requests_aws4auth import AWS4Auth import re from datetime import datetime # 初始化客户端 s3 = boto3.client('s3') region = 'us-east-1' # 替换为你的AWS区域 service = 'es' credentials = boto3.Session().get_credentials() awsauth = AWS4Auth(credentials.access_key, credentials.secret_key, region, service, session_token=credentials.token) # OpenSearch客户端配置 host = 'your-opensearch-endpoint' # 替换为你的OpenSearch端点 os_client = OpenSearch( hosts = [{'host': host, 'port': 443}], http_auth = awsauth, use_ssl = True, verify_certs = True, connection_class = RequestsHttpConnection ) # Apache日志解析正则(根据你的实际日志格式调整) APACHE_LOG_PATTERN = r'^(\S+) (\S+) (\S+) \[([\w:/]+\s[+\-]\d{4})\] "(\S+) (\S+) (\S+)" (\d{3}) (\d+)$' def lambda_handler(event, context): for record in event['Records']: # 解析SQS消息体为JSON字典 try: s3_event = json.loads(record['body']) except json.JSONDecodeError as e: print(f"解析SQS消息体失败: {e}") continue # 提取S3桶和对象键(解码URL编码) bucket = s3_event['Records'][0]['s3']['bucket']['name'] key = urllib.parse.unquote_plus(s3_event['Records'][0]['s3']['object']['key'], encoding='utf-8') # 获取S3对象内容 try: obj = s3.get_object(Bucket=bucket, Key=key) log_content = obj['Body'].read().decode('utf-8') log_lines = log_content.splitlines() except Exception as e: print(f"读取S3对象失败: {e}") continue # 解析日志并批量准备OpenSearch文档 bulk_docs = [] for line in log_lines: match = re.match(APACHE_LOG_PATTERN, line) if not match: print(f"无法解析日志行: {line}") continue # 提取日志字段 remote_host, remote_logname, remote_user, timestamp_str, method, path, protocol, status, size = match.groups() # 转换时间格式为ISO8601 timestamp = datetime.strptime(timestamp_str, '%d/%b/%Y:%H:%M:%S %z').isoformat() # 构造OpenSearch文档 doc = { 'remote_host': remote_host, 'remote_user': remote_user, 'timestamp': timestamp, 'method': method, 'path': path, 'protocol': protocol, 'status': int(status), 'size': int(size), 'source_file': key } # 批量操作指令 bulk_docs.append({'index': {'_index': 'apache-access-logs'}}) bulk_docs.append(doc) # 批量写入OpenSearch if bulk_docs: try: response = os_client.bulk(body=bulk_docs) if response['errors']: print(f"OpenSearch批量写入错误: {response['items']}") else: print(f"成功写入{len(bulk_docs)//2}条日志到OpenSearch") except Exception as e: print(f"OpenSearch写入失败: {e}")
关键说明
- JSON解析:必须用
json.loads()将SQS的字符串消息体转换为Python字典,才能进行键值访问 - URL解码:S3对象键中的特殊字符(如空格、中文)会被自动编码,需用
urllib.parse.unquote_plus还原 - 日志解析:正则表达式需匹配你的Apache日志实际格式,常见格式参考NCSA组合日志格式
- OpenAuth认证:使用
AWS4Auth实现Lambda到OpenSearch的IAM身份验证,无需硬编码密钥
内容的提问来源于stack exchange,提问作者shicha
相关产品推荐
相关产品推荐

