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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 20:22:19