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

SQS触发Lambda读取S3对象遇TypeError报错求助

解决SQS触发Lambda读取S3日志并导入OpenSearch的问题

问题根源

你碰到的TypeError: string indices must be integers,核心原因是从SQS消息中提取的record['body']是JSON格式字符串,而非Python字典,直接用字典索引访问['Records']会报错。你的代码仅做了赋值操作,没有完成字符串到JSON对象的解析。

解决方案步骤

  1. 导入json模块,将SQS消息体的JSON字符串解析为Python字典
  2. 处理S3对象键的URL编码(比如空格、特殊字符会被SQS转义,需解码还原)
  3. 读取S3中的Apache日志文件,逐行解析日志格式
  4. 将解析后的日志批量写入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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 10:45:37