AWS SNS发送Protocol Buffer消息至FastAPI订阅端解码失败排查
SNS二进制消息属性解码错误问题
消息已成功发送,但订阅端接收时抛出(Decode Error) Error parsing message解码错误。原本以为发送的是Binary类型消息无需额外解码,但实际收到的是字符串,必须调整解码步骤。
发布端代码(publisher.py)
import boto3 from my_generated_pb2 import MyMessage sns_client = boto3.client("sns") message = MyMessage() message.content = "Hello, SNS!" response = sns_client.publish( TopicArn="TOPIC_ARN", Message='test', MessageAttributes={ 'ProtocolBuffer': { 'DataType': 'Binary', 'BinaryValue': message.SerializeToString() }, } )
订阅端代码(subscriber.py)
from fastapi import FastAPI from my_generated_pb2 import MyMessage app = FastAPI() @app.post("/subscriber") async def subscriber(request: Request): """ Receives a message from SNS. :param request: :return: """ body = await request.json() if body['Type'] == "Notification": serialized_message = body['MessageAttributes']['ProtocolBuffer']['Value'].encode() message_obj = TransactionRequest() message_obj.ParseFromString(serialized_message) # ERROR!
本地测试正常代码
from my_generated_pb2 import MyMessage message = MyMessage() message.content = "Hello, SNS!" serialized_message = message.SerializeToString() message_obj = MyMessage() message_obj.ParseFromString(serialized_message)
问题原因与修复方案
这不是Boto或SNS的bug,是对SNS二进制消息属性传输机制的误解:
- SNS会自动对二进制类型的消息属性进行Base64编码后传输,订阅端收到的
Value是Base64编码字符串,而非原始二进制数据。 - 直接调用
encode()仅能将字符串转为UTF-8字节流,无法还原Protobuf序列化的原始二进制数据,因此触发解码失败。
修复后的订阅端代码
from fastapi import FastAPI, Request import base64 from my_generated_pb2 import MyMessage app = FastAPI() @app.post("/subscriber") async def subscriber(request: Request): body = await request.json() if body['Type'] == "Notification": # 对Base64编码字符串解码,得到原始二进制数据 base64_str = body['MessageAttributes']['ProtocolBuffer']['Value'] serialized_message = base64.b64decode(base64_str) message_obj = MyMessage() message_obj.ParseFromString(serialized_message) # 可正常解码 return {"status": "success", "content": message_obj.content}
同时注意:原订阅端代码中实例化的TransactionRequest()与导入的MyMessage不一致,需保持类型名称统一,避免额外错误。
内容的提问来源于stack exchange,提问作者J. Javier Gálvez
相关产品推荐
相关产品推荐

