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

如何在Python中提取SQS JSON消息嵌套值?AWS链路报错排查

问题:无法从SQS消息中提取vote字段值

链路架构

  • Lambda -> SNS -> SQS -> EC2上的Python脚本

测试事件结构(已启用原始消息投递)

{  
  "body": {    
    "MessageAttributes": {      
      "vote": {        
        "Type": "Number",        
        "Value": "90"      
      },      
      "voter": {        
        "Type": "String",        
        "Value": "default_voter"      
      }    
    }  
  }

当前Python处理脚本

#!/usr/bin/env python3

import boto3
import json
import logging
import sys

logging.basicConfig(stream=sys.stdout, level=logging.INFO)

queue = boto3.resource('sqs', region_name='us-east-1').get_queue_by_name(QueueName="erjan")
table = boto3.resource('dynamodb', region_name='us-east-1').Table('Votes')

def process_message(message):
    try:
        payload = json.loads(message) #unable to parse sqs json msg here
        #payload = message.body
        #payload = message['body']
        voter = payload['MessageAttributes']['voter']['Value'] #here the exception raised!
        vote  = payload['MessageAttributes']['vote']['Value']
        logging.info("Voter: %s, Vote: %s", voter, vote)
        store_vote(voter, vote)
        update_count(vote)
        message.delete()
    except Exception as e:
        print('x = msg.body')
        x = (message.body)
        print(x)
        print('-------')
        print('message.body')
        print(message.body)

        try:
            vote = x['MessageAttributes']['vote']['Value']
            logging.error("Failed to process message")
            logging.error('------- here: ' + str(e))
            logging.error('vote %d' % vote)
        except TypeError:
            logging.error("error catched")

def store_vote(voter, vote):
    try:
        logging.info('table put item.......')
        print('table put item......')
        response = table.put_item(
           Item={'voter': voter, 'vote': vote}
        )
    except:
        logging.error("Failed to store message")
        raise

def update_count(vote):
    logging.info('update count....')
    print('update count....')
    table.update_item(
        Key={'voter': 'count'},
        UpdateExpression="set #vote = #vote + :incr",
            ExpressionAttributeNames={'#vote': vote},
            ExpressionAttributeValues={':incr': 1}
    )

if __name__ == "__main__":
    while True:
        try:
            messages = queue.receive_messages()
        except KeyboardInterrupt:
           logging.info("Stopping...")
           break
        except:
            logging.error(sys.exc_info()[0])
            logging.info('here error - we continue')
            continue
        for message in messages:
            process_message(message)

报错信息

payload = json.loads(message)
  File "/usr/lib64/python3.7/json/__init__.py", line 341, in loads
    raise TypeError(f'the JSON object must be str, bytes or bytearray, '
TypeError: the JSON object must be str, bytes or bytearray, not sqs.Message

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "./processor.py", line 90, in <module>
    process_message(message)
  File "./processor.py", line 40, in process_message
    vote = x['MessageAttributes']['vote']['Value']
TypeError: string indices must be integers

问题现象

  • 代码中误将sqs.Message对象传入json.loads,而非消息体字符串message.body
  • 打印message仅输出SQS消息元数据,无法直接查看消息内容
  • 异常处理中直接对字符串格式的消息体使用字典索引,触发类型错误

解决方案

  1. 修复核心解析逻辑:将json.loads(message)改为json.loads(message.body),json.loads仅接受字符串/字节类型参数,message.body才是存储实际消息内容的字段。
  2. 匹配测试事件的嵌套结构:测试事件中MessageAttributes在body字段内部,需先解析外层JSON,再提取内层内容:
def process_message(message):
    try:
        # 先解析外层JSON
        payload = json.loads(message.body)
        # 提取body内部的消息属性
        body_content = payload['body']
        voter = body_content['MessageAttributes']['voter']['Value']
        vote = body_content['MessageAttributes']['vote']['Value']
        
        logging.info("Voter: %s, Vote: %s", voter, vote)
        store_vote(voter, vote)
        update_count(vote)
        message.delete()
    except Exception as e:
        logging.error(f"Processing failed: {str(e)}")
        # 异常处理中也需先解析字符串为字典
        try:
            raw_body = message.body
            if raw_body:
                payload = json.loads(raw_body)
                vote = payload['body']['MessageAttributes']['vote']['Value']
                logging.error(f"Fallback extracted vote: {vote}")
        except Exception as fallback_err:
            logging.error(f"Fallback parsing failed: {str(fallback_err)}")
  1. 验证消息完整性:确认SNS原始消息投递功能正常,Lambda发送到SNS的消息结构与测试事件一致,确保SQS能收到完整的JSON格式消息体。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 05:30:47