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

