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

Python如何解析Lambda接收的Athena任务SQS队列消息体并提取所需字段

解决方案

核心问题梳理

你之前的实现存在两个明显问题:

  • 初始状态下Step Function输出的Athena结果没有给key和字符串值加双引号,属于非标准JSON,无法直接用json库解析,这个问题你已经通过修正Athena输出配置解决了。
  • 你写的测试代码逻辑有误:把原始事件字典用json.dumps转成了字符串,后续直接遍历字符串肯定会报错。

正确实现代码

import json

# 注意:Lambda触发时收到的event本身就是Python字典,不需要手动做json.dumps处理
# 这里模拟测试直接定义事件字典即可
test_event = {
    'Records': [
        {
            'messageId': '1ca8da82-e973-4c80-9a62-d8515fe8e436',
            'body': '{"ResultSet":{"ResultSetMetadata":{"ColumnInfo":[{"CaseSensitive":false,"CatalogName":"hive","Label":"eventname","Name":"eventname","Nullable":"UNKNOWN","Precision":0,"Scale":0,"SchemaName":"","TableName":"","Type":"json"},{"CaseSensitive":false,"CatalogName":"hive","Label":"eventsource","Name":"eventsource","Nullable":"UNKNOWN","Precision":0,"Scale":0,"SchemaName":"","TableName":"","Type":"json"},{"CaseSensitive":false,"CatalogName":"hive","Label":"awsregion","Name":"awsregion","Nullable":"UNKNOWN","Precision":0,"Scale":0,"SchemaName":"","TableName":"","Type":"json"},{"CaseSensitive":false,"CatalogName":"hive","Label":"username","Name":"username","Nullable":"UNKNOWN","Precision":0,"Scale":0,"SchemaName":"","TableName":"","Type":"json"}]},"Rows":[{"Data":[{"VarCharValue":"eventname"},{"VarCharValue":"eventsource"},{"VarCharValue":"awsregion"},{"VarCharValue":"username"}]},{"Data":[{"VarCharValue":"\"DeleteBucket\""},{"VarCharValue":"\"s3.amazonaws.com\""},{"VarCharValue":"\"us-west-2\""},{"VarCharValue":"\"AROA4SALNAMTBCVMSUKMJ:travis.jorge@tylerhost.net\""}]},{"Data":[{"VarCharValue":"\"StartExecution\""},{"VarCharValue":"\"states.amazonaws.com\""},{"VarCharValue":"\"us-west-2\""},{"VarCharValue":"\"AROA4SALNAMTBCVMSUKMJ:travis.jorge@tylerhost.net\""}]},{"Data":[{"VarCharValue":"\"BatchDeleteTable\""},{"VarCharValue":"\"glue.amazonaws.com\""},{"VarCharValue":"\"us-west-2\""},{"VarCharValue":"\"AROA4SALNAMTBCVMSUKMJ:travis.jorge@tylerhost.net\""}]}]},"UpdateCount":0}',
            # 其他字段省略不影响逻辑
        }
    ]
}

def lambda_handler(event, context):
    for record in event['Records']:
        # 1. 把SQS消息体的JSON字符串转为Python字典
        athena_result = json.loads(record['body'])
        result_set = athena_result['ResultSet']
        # 2. Rows数组第一个元素是表头,跳过从下标1开始取实际数据
        for row in result_set['Rows'][1:]:
            data = row['Data']
            # 3. 提取字段,同时去掉VarCharValue外层多余的双引号
            event_type = data[0]['VarCharValue'].strip('"')
            region = data[2]['VarCharValue'].strip('"')
            user_name = data[3]['VarCharValue'].strip('"')
            # 后续业务逻辑
            print(f"事件类型:{event_type},区域:{region},用户名:{user_name}")

# 本地测试调用
if __name__ == "__main__":
    lambda_handler(test_event, None)

更优方案建议

如果可以调整上游Step Function的逻辑,推荐做如下优化:

  • 不要把整个Athena结果集直接塞到SQS消息里,在Step Function中先提取你需要的event_type、region、user_name字段,按业务需要拼装成精简结构后再发SQS,既减少消息体积,也避免Lambda侧做冗余解析。
  • 如果Athena返回的结果行数很多,结果集大小超过SQS单条消息256KB的限制,可以让Step Function只把Athena结果的S3存储路径发到SQS,Lambda收到消息后去S3下载完整结果文件再处理,避免消息溢出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 14:24:04