使用 Watermill 在 Go 中集成 AWS SNS/SQS:从队列选型到本地模拟器联调
【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill
AWS SQS(Simple Queue Service)与 SNS(Simple Notification Service)是亚马逊云提供的全托管消息服务,用于解耦并弹性扩展微服务、分布式系统与 Serverless 应用。Watermill 通过独立的watermill-aws包封装了 AWS SDK v2 的全部底层细节,向 Go 开发者暴露与 Watermill 其它 Pub/Sub 完全一致的消息发布/订阅 API。读完本文,你将掌握 SQS/SNS 的选型逻辑、所需的 IAM 权限、Subscriber/Publisher 配置、队列 URL 与 Topic ARN 的解析机制,并能基于 LocalStack 在本地跑通官方示例做端到端验证。
概述:Watermill 如何对接 AWS 消息服务
在 Watermill 的模型里,一切消息传递都围绕message.Publisher与message.Subscriber两个接口展开(见 message/pubsub.go):
Publisher.Publish(topic string, messages ...*Message) error负责向指定 topic 发布消息;Subscriber.Subscribe(ctx context.Context, topic string) (<-chan *Message, error)返回一个消息通道,处理完成后必须调用msg.Ack(),失败需要重投则调用msg.Nack();- 当订阅端
ctx被取消时,订阅关闭且输出通道随之关闭。
watermill-aws包正是把 AWS SQS 的队列与 SNS 的 Topic 归一化到这一对接口之上:它对调用方屏蔽了 AWS SDK 的初始化、长轮询、消息可见性超时等内部机制,你只需要提供AWSConfig和若干可选配置即可完成发布与订阅。
安装 watermill-aws
go get github.com/ThreeDotsLabs/watermill-aws仓库中的官方示例(见 _examples/pubsubs/aws-sqs/go.mod 与 _examples/pubsubs/aws-sns/go.mod)当前依赖github.com/ThreeDotsLabs/watermill v1.5.1、github.com/ThreeDotsLabs/watermill-aws v1.0.2,并基于 AWS SDK v2(github.com/aws/aws-sdk-go-v2)实现,同时使用github.com/samber/lo作为辅助库。这意味着watermill-aws面向的是 aws-sdk-go-v2 的新一代 API,而非旧版aws-sdk-go。
SQS 与 SNS:先选型,再编码
SQS 与 SNS 虽然同属 AWS 消息服务,但定位完全不同,适用场景也截然不同,动手写代码前必须先明确需求。
SNS 与 SQS 是如何连接起来的
要把 SNS 当作 Pub/Sub 使用(即让多个订阅者收到同一条消息),标准做法是:创建一个 SNS Topic,然后让多个 SQS 队列订阅它。消息发布到 SNS Topic 后,会被投递给所有已订阅的 SQS 队列。watermill-aws已经内置了这套逻辑——当你调用Subscribe订阅某个 SNS Topic 时,Watermill 会自动创建一个 SQS 队列并完成对该 Topic 的订阅。
因此可以认为:单个 SQS 队列在语义上等价于其它 Pub/Sub 实现中的一个消费组(consumer group)或一个订阅(subscription)。这条连接机制与 AWS 官方文档中"将 SQS 队列订阅到 SNS Topic"的说明完全一致。
SQS(Simple Queue Service)
- 适用场景:需要一个面向单一消费者的简单消息队列;
- 擅长:任务队列、后台作业处理;
- 能力:通过 FIFO 队列支持 exactly-once 处理,并(在多数情况下)保证消息顺序;
- 典型用例:后台处理用户上传的文件。
SNS(Simple Notification Service)
- 适用场景:需要把消息广播给多个订阅者;
- 擅长:实现 pub/sub 模式、支撑事件驱动架构;
- 能力:支持多种订阅端类型(SQS、Lambda、HTTP/S、Email、SMS 等);
- 典型用例:新用户注册后通知多个下游服务;
- Watermill 的 SNS 实现会自动为每个订阅者创建并管理对应的 SQS 队列,大幅简化"一个 Topic 对接多个 SQS 队列"的运维工作。
两者并非互斥:你完全可以在同一个应用里同时使用——用 SNS 广播领域事件,用 SQS 处理由这些事件触发的具体任务。
使用 SQS:队列即 Topic
特性一览
| 特性 | 实现情况 | 说明 |
|---|---|---|
| ConsumerGroups | 否 | SQS 本质是队列,如需消费组能力请使用 SNS |
| ExactlyOnceDelivery | 否 | AWS 仅在 FIFO 队列上提供 exactly-once 处理 |
| GuaranteedOrder | 是* | 标准队列因高度分布式架构可能重复投递或偶发乱序,AWS 会尽力维持发送顺序 |
| Persistent | 是 | 消息持久化存储 |
其中 "GuaranteedOrder" 一栏的星号值得注意:AWS 官方对标准队列的描述明确写道,"由于高度分布式的架构,同一条消息可能有多个副本被投递,消息偶尔会乱序到达。尽管如此,标准队列仍会尽力维持消息的发送顺序"。因此 SQS 的"顺序保证"是 best-effort 级别的。
所需的 IAM 权限
要让 Watermill 的 SQS 发布/订阅正常工作,IAM 策略至少需要授予:
sqs:ReceiveMessagesqs:DeleteMessagesqs:GetQueueUrlsqs:CreateQueuesqs:GetQueueAttributessqs:SendMessagesqs:ChangeMessageVisibility
SQS 配置结构
SQS 的SubscriberConfig与PublisherConfig定义在watermill-aws包的sqs/config.go中,核心字段包括AWSConfig(AWS SDK v2 的全局配置,如凭证、Region)与OptFns(一组func(*amazonsqs.Options)回调,用于透传并覆盖 AWS SDK 的底层选项,例如自定义 Endpoint)。从官方示例 main.go 可以看到最基本的组装方式:
subscriberConfig := sqs.SubscriberConfig{ AWSConfig: aws.Config{ Credentials: aws.AnonymousCredentials{}, }, OptFns: sqsOpts, } subscriber, err := sqs.NewSubscriber(subscriberConfig, logger) if err != nil { panic(err) } messages, err := subscriber.Subscribe(context.Background(), "example-topic")发布端结构完全对称:
publisherConfig := sqs.PublisherConfig{ AWSConfig: aws.Config{ Credentials: aws.AnonymousCredentials{}, }, OptFns: sqsOpts, } publisher, err := sqs.NewPublisher(publisherConfig, logger) if err != nil { panic(err) }解析队列 URL:把 AWS 资源归一化为 topic
在 Watermill 的模型中,发布与订阅方法接收的是一个字符串topic。对 SQS 而言,watermill-aws会把AWS 队列 URL 归一化为这个 topic,让你不必在业务代码里处处拼写完整的队列 URL。
为此,包内定义了QueueUrlResolver接口,允许你自定义"Watermill topic → AWS 队列 URL"的解析逻辑。默认使用GetQueueUrlByNameUrlResolver——它按队列名向 AWS 查询并解析出真实队列 URL。此外还提供了两个内置解析器:
GenerateQueueUrlResolver:根据命名规则直接生成队列 URL,无需额外的 AWS API 调用;TransparentUrlResolver:把传入的 topic 原样当作队列 URL 使用,不做任何转换。
如果默认行为不满足需求,实现自己的QueueUrlResolver并在SubscriberConfig/PublisherConfig中注入即可。
在本地用 SQS 模拟器联调
本地开发或测试时,可以选用goaws或localstack模拟 AWS 服务,并通过OptFns覆盖 Endpoint:
package main import ( amazonsqs "github.com/aws/aws-sdk-go-v2/service/sqs" "github.com/ThreeDotsLabs/watermill-amazonsqs/sqs" ) func main() { // ... sqsOpts := []func(*amazonsqs.Options){ amazonsqs.WithEndpointResolverV2(sqs.OverrideEndpointResolver{ Endpoint: transport.Endpoint{ URI: *lo.Must(url.Parse("http://localstack:4566")), }, }), } sqsConfig := sqs.SubscriberConfig{ AWSConfig: cfg, OptFns: sqsOpts, } sub, err := sqs.NewSubscriber(sqsConfig, logger) if err != nil { panic(fmt.Errorf("unable to create new subscriber: %w", err)) } // ... }仓库中的 docker-compose.yml 给出了完整的本地环境:使用localstack/localstack:3.8镜像,通过环境变量SERVICES=sqs,sns同时启用 SQS 与 SNS,默认 Region 为us-east-1,对外暴露4566-4597端口,并配置了健康检查(awslocal sqs list-queues)。配合 Go 1.25 的golang:1.25容器执行go run main.go,即可在不产生任何真实 AWS 费用的情况下跑通完整链路。
使用 SNS:多订阅者广播
特性一览
| 特性 | 实现情况 | 说明 |
|---|---|---|
| ConsumerGroups | 是 | 每个订阅者对应一个自动创建的 SQS 队列 |
| ExactlyOnceDelivery | 否 | AWS 仅在 FIFO 队列上提供 exactly-once 处理 |
| GuaranteedOrder | 是* | 底层依赖 SQS 标准队列,顺序保证为 best-effort |
| Persistent | 是 | 消息持久化存储 |
所需的 IAM 权限
SNS 侧需要:
sns:Subscribesns:ConfirmSubscriptionsns:Receivesns:Unsubscribe
同时,因为 SNS 订阅者底层就是 SQS 队列,还需要完整的 SQS 权限集:
sqs:ReceiveMessagesqs:DeleteMessagesqs:GetQueueUrlsqs:CreateQueuesqs:GetQueueAttributessqs:SendMessagesqs:ChangeMessageVisibilitysqs:SetQueueAttributes
补充说明:当sns.SubscriberConfig.DoNotSetQueueAccessPolicy未开启时,Watermill 会自动为自动创建的队列设置访问策略,以便 SNS 向该队列投递消息,因此此时还额外需要sqs:SetQueueAttributes权限。
SNS 配置结构:Subscriber 需要双份配置
SNS 的SubscriberConfig定义在watermill-aws包的sns/config.go中,同样包含AWSConfig、OptFns等字段,并且因为 SNS 订阅者使用 SQS 队列作为"订阅",构造sns.NewSubscriber时还必须同时传入一份 SQS 配置。这一签名在官方示例 main.go 中体现得很直观:
return sns.NewSubscriber(subscriberConfig, sqsSubscriberConfig, logger)SNS 的SubscriberConfig还支持两个关键字段:
TopicResolver:把 Watermill topic 解析为 AWS Topic ARN;GenerateSqsQueueName:func(ctx context.Context, snsTopic sns.TopicArn) (string, error),为每个订阅者生成其专属 SQS 队列名。
官方 SNS 示例正是利用GenerateSqsQueueName实现"每个订阅者一个队列"的消费组语义:它先用sns.ExtractTopicNameFromTopicArn从 ARN 中取出 Topic 名,再拼接订阅者后缀,于是订阅者 A 与 B 分别得到形如example-topic-subA与example-topic-subB的独立队列;同一 Topic 上两个订阅者会各自收到同一条消息的完整副本,这正是 SNS 广播能力的体现:
newSubscriber := func(name string) (message.Subscriber, error) { subscriberConfig := sns.SubscriberConfig{ AWSConfig: aws.Config{ Credentials: aws.AnonymousCredentials{}, }, OptFns: snsOpts, TopicResolver: topicResolver, GenerateSqsQueueName: func(ctx context.Context, snsTopic sns.TopicArn) (string, error) { topic, err := sns.ExtractTopicNameFromTopicArn(snsTopic) if err != nil { return "", err } return fmt.Sprintf("%v-%v", topic, name), nil }, } sqsSubscriberConfig := sqs.SubscriberConfig{ AWSConfig: aws.Config{ Credentials: aws.AnonymousCredentials{}, }, OptFns: sqsOpts, } return sns.NewSubscriber(subscriberConfig, sqsSubscriberConfig, logger) }解析 Topic ARN:把 AWS 资源归一化为 topic
与 SQS 把队列 URL 归一化为 topic 类似,SNS 侧watermill-aws会把AWS Topic ARN 归一化为 Watermill topic,由TopicResolver接口完成,包内提供两个开箱即用的解析器:
TransparentTopicResolver:把传入的 topic 原样当作 ARN 使用;GenerateArnTopicResolver:根据账号 ID 与 Region 直接生成 ARN,无需查询 AWS。
GenerateArnTopicResolver在本地模拟场景下尤其常用。由于 LocalStack 并不提供真实的账号与区域环境,示例代码使用sns.NewGenerateArnTopicResolver("000000000000", "us-east-1")显式构造解析器,使Publish/Subscribe中传入的"example-topic"能被映射为arn:aws:sns:us-east-1:000000000000:example-topic形式的真实 ARN:
topicResolver, err := sns.NewGenerateArnTopicResolver("000000000000", "us-east-1") if err != nil { panic(err) }在本地用 SNS 模拟器联调
与 SQS 一样,可通过OptFns覆盖 Endpoint 指向 localstack 或 goaws:
package main import ( amazonsns "github.com/aws/aws-sdk-go-v2/service/sns" "github.com/ThreeDotsLabs/watermill-amazonsns/sns" ) func main() { // ... snsOpts := []func(*amazonsns.Options){ amazonsns.WithEndpointResolverV2(sns.OverrideEndpointResolver{ Endpoint: transport.Endpoint{ URI: *lo.Must(url.Parse("http://localstack:4566")), }, }), } snsConfig := sns.SubscriberConfig{ AWSConfig: cfg, OptFns: snsOpts, } sub, err := sns.NewSubscriber(snsConfig, sqsConfig, logger) if err != nil { panic(fmt.Errorf("unable to create new subscriber: %w", err)) } // ... }注意这里sns.NewSubscriber接收的是(snsConfig, sqsConfig, logger)三个参数——再次印证 SNS 订阅依赖 SQS 队列这一事实。
端到端示例解读:Ack 语义与消息处理
无论 SQS 还是 SNS,官方示例都展示了 Watermill 统一的消息消费范式(见 aws-sqs/main.go 与 aws-sns/main.go):
func publishMessages(publisher message.Publisher) { for { msg := message.NewMessage(watermill.NewUUID(), []byte("Hello, world!")) if err := publisher.Publish("example-topic", msg); err != nil { panic(err) } time.Sleep(time.Second) } } func process(messages <-chan *message.Message) { for msg := range messages { log.Printf("received message: %s, payload: %s", msg.UUID, string(msg.Payload)) // 必须调用 Ack 确认已收到并处理完成, // 否则该消息会被反复重投。 msg.Ack() } }几个要点:
- 消息 ID 使用
watermill.NewUUID()生成,消费者通过msg.UUID与msg.Payload读取消息; - 必须调用
msg.Ack()——这与 message/pubsub.go 中Subscribe的契约一致:只有 Ack 后才会拉取下一条消息;处理失败应改调Nack()触发重投。若既不 Ack 也不 Nack,消息会因可见性超时而反复重投; - 在 SNS 示例中,同一条消息会被
subA与subB两个订阅者各自完整收到一次(日志中前缀 A/B 分别打印),直观演示了"一个 Topic 广播到多个队列"的效果。
本地运行时,直接使用示例目录下的 docker-compose.yml(SERVICES=sqs,sns,LocalStack 3.8)即可一键启动模拟环境并观察上述行为,全程无需真实 AWS 账号。
小结
watermill-aws把 AWS 的两类消息服务收敛进了 Watermill 统一的Publisher/Subscriber模型:SQS 适合单消费者队列与任务处理,SNS 适合多订阅者广播与事件驱动架构;SNS 订阅底层自动创建 SQS 队列,使单个队列等价于一个消费组。实际使用时,只需关心 IAM 权限、AWSConfig/OptFns配置以及 URL/ARN 解析器的选择,其余 SDK 细节均由包内封装。配合官方示例与 LocalStack 模拟器,你可以在完全本地化的环境中完成开发、测试与验证。
【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考