拓冰建站拓冰建站
首页 / 资讯中心 / 正文

在 Watermill 中使用 Google Cloud Pub/Sub:安装、配置与实战指南

在 Watermill 中使用 Google Cloud Pub/Sub安装、配置与实战指南【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermillGoogle Cloud Pub/Sub 是 Google Cloud Platform 提供的全托管实时消息服务将企业级消息中间件的灵活性与可靠性带到了云端。本文以 Watermill 项目官方文档 docs/content/pubsubs/googlecloud.md 为主线完整讲解如何通过watermill-googlecloud/v2适配器在 Go 应用中发布与订阅消息从安装、环境变量配置、模拟器本地开发到订阅命名、Marshaler 定制与完整可运行示例读完即可在真实 GCP 环境或本地模拟器中跑通第一条消息链路。Google Cloud Pub/Sub 与 Watermill 的集成概览Cloud Pub/Sub 是一个可扩展、持久化的事件摄取与投递系统是现代流分析管线的底层基础设施。它提供多对多、异步的消息传递将发送方与接收方解耦使得独立编写的应用之间可以进行安全、高可用的通信同时以低延迟、持久化的消息投递帮助开发者快速集成 GCP 内外的各类系统。在 Watermill 的生态中Google Cloud Pub/Sub 由独立的适配器模块github.com/ThreeDotsLabs/watermill-googlecloud/v2提供支持参见 README.md 的 Pub/Subs 一节。它实现了 Watermill 定义的标准message.Publisher与message.Subscriber接口定义见 message/pubsub.gotype Publisher interface { Publish(topic string, messages ...*Message) error Close() error } type Subscriber interface { Subscribe(ctx context.Context, topic string) (-chan *Message, error) Close() error }这意味着只要遵循这套接口应用代码无需感知底层是 Kafka、RabbitMQ 还是 Google Cloud Pub/Sub业务逻辑与消息基础设施完全解耦。特性一览在选用任何 Pub/Sub 之前都需要先确认它对关键消息语义的支持程度。Google Cloud Pub/Sub 在 Watermill 中的支持情况如下引自官方文档FeatureImplementsNoteConsumerGroupsyes同一 Subscription 名称下的多个订阅者多个消费者实例共享消费ExactlyOnceDeliveryno不提供恰好一次投递GuaranteedOrderno不保证消息顺序Persistentyes*消息最大保留时间为 7 天其中Persistent一栏的*号需要特别注意Google Cloud Pub/Sub 的消息是持久化的但最大保留时间是 7 天超出保留期的消息将被系统清理这一点在设计事件回放或长周期补偿场景时需要纳入考虑。安装Google Cloud Pub/Sub 适配器是一个独立于 Watermill 主仓库的 Go module。安装方式与安装其他 Go 依赖一致go get github.com/ThreeDotsLabs/watermill-googlecloud/v2当前仓库中的示例模块使用了v2.0.0版本见 _examples/pubsubs/googlecloud/go.mod如果你使用的是go mod管理模式上述命令会自动将适配器及其传递依赖写入go.sum。安装后在代码中通过以下路径导入import github.com/ThreeDotsLabs/watermill-googlecloud/v2/pkg/googlecloud连接与认证配置Watermill 会通过环境变量连接到 Google Cloud Pub/Sub 实例因此你的应用不需要硬编码任何连接串。生产环境GOOGLE_APPLICATION_CREDENTIALS在生产环境需要设置GOOGLE_APPLICATION_CREDENTIALS环境变量指向包含服务账号凭据的 JSON 文件路径。这是 Google Cloud 客户端库的标准认证方式。一个关键优势是你不需要安装 Google Cloud SDKgcloud。Watermill 会替你完成管理类任务——例如在默认配置和权限允许的情况下自动创建 topic 和 subscription。开发环境Pub/Sub 模拟器对于本地开发Google 官方提供了 Pub/Sub 模拟器emulator配合PUBSUB_EMULATOR_HOST环境变量即可完全脱离 GCP 云端进行开发调试。仓库中的 _examples/pubsubs/googlecloud/docker-compose.yml 给出了开箱即用的模拟器编排services: server: image: golang:1.25 restart: unless-stopped depends_on: - googlecloud volumes: - .:/app - $GOPATH/pkg/mod:/go/pkg/mod environment: # use local emulator instead of google cloud engine PUBSUB_EMULATOR_HOST: googlecloud:8085 working_dir: /app command: go run main.go googlecloud: image: google/cloud-sdk:414.0.0 entrypoint: gcloud --quiet beta emulators pubsub start --host-port0.0.0.0:8085 --verbositydebug --log-http restart: unless-stopped这个编排做了两件事googlecloud服务使用官方google/cloud-sdk镜像启动 Pub/Sub 模拟器监听0.0.0.0:8085server服务挂载当前目录作为工作区并执行go run main.go同时通过PUBSUB_EMULATOR_HOST: googlecloud:8085将应用指向模拟器。在开发环境运行docker-compose up即可一键拉起整套环境更详细的运行方式可参考 docs/content/learn/getting-started.md 中的 Running in Docker 一节。Publisher 与 Subscriber 配置适配器的两个核心配置结构体是PublisherConfig与SubscriberConfig二者均位于pkg/googlecloud包内。结合当前仓库中可验证的示例代码_examples/pubsubs/googlecloud/main.go最基本的使用方式是publisher, err : googlecloud.NewPublisher(googlecloud.PublisherConfig{ ProjectID: test-project, }, logger) if err ! nil { panic(err) } subscriber, err : googlecloud.NewSubscriber( googlecloud.SubscriberConfig{ // custom function to generate Subscription Name, // there are also predefined TopicSubscriptionName and TopicSubscriptionNameWithSuffix available. GenerateSubscriptionName: func(topic string) string { return test-sub_ topic }, ProjectID: test-project, }, logger, ) if err ! nil { panic(err) }其中ProjectID指定消息所属的 GCP 项目 ID无论发布还是订阅都必须填写SubscriberConfig.GenerateSubscriptionName是一个接收 topic 名、返回订阅名的函数用于控制订阅的命名策略详见下文订阅命名一节NewPublisher/NewSubscriber的第二个参数是 Watermill 的 logger示例中使用watermill.NewStdLogger(false, false)。关于配置结构体的完整字段以watermill-googlecloud模块源码中PublisherConfig/SubscriberConfig结构体的实际定义为准官方文档通过load-snippet-partial直接内嵌了这两段源码。发布消息发布操作通过Publisher.Publish完成。其底层签名官方文档内嵌源码为func (p *Publisher) Publish(topic string, messages ...*message.Message) errorWatermill 的message.Message由 UUID、Payload 与 Metadata 组成示例中通过message.NewMessage(watermill.NewUUID(), []byte(Hello, world!))构造消息再发布到指定 topicfunc 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) } }注意Publish接收message.Publisher接口而非具体的*googlecloud.Publisher这也是 Watermill 解耦思想的体现——调用方只依赖接口。订阅消息订阅操作通过Subscriber.Subscribe完成。其底层签名官方文档内嵌源码为func (s *Subscriber) Subscribe(ctx context.Context, topic string) (-chan *message.Message, error)重要语义Subscribe调用时会自动创建订阅如果需要并且只有订阅创建之后发布到 topic 的消息才会被接收。这一点与 Kafka 等消费组模型不同是使用 Google Cloud Pub/Sub 时必须牢记的心智模型。订阅返回一个只读的-chan *message.Message通道消费方从通道中不断取出消息处理func process(messages -chan *message.Message) { for msg : range messages { log.Printf(received message: %s, payload: %s, msg.UUID, string(msg.Payload)) // we need to Acknowledge that we received and processed the message, // otherwise, it will be resent over and over again. msg.Ack() } }这里msg.Ack()是必须调用的如果不对消息进行确认acknowledgePub/Sub 会认为消息未被处理成功从而反复重新投递。所有消费逻辑完成后务必确认消息否则可能造成消息风暴。深入理解 Subscription name在 Google Cloud Pub/Sub 中topic 与 subscription 是两个独立的资源理解二者的关系是正确使用的前提要接收发布到某个 topic 的消息必须先为该 topic 创建订阅subscription只有订阅创建之后发布到该 topic 的消息才对订阅方应用可见订阅将 topic 连接到接收并处理消息的订阅方应用一个 topic 可以有多个订阅但一个订阅只能属于一个 topic。在 Watermill 中订阅是在调用Subscribe()时自动创建的订阅名称由传入SubscriberConfig.GenerateSubscriptionName的函数生成。默认情况下该函数是TopicSubscriptionName——也就是直接使用 topic 名作为订阅名。多消费者场景下的命名策略当一个 topic 需要被多个消费者实例消费时即文档特性表中的 ConsumerGroups yes需要注意使用相同的订阅名多个订阅者会作为同一个消费组协同消费——每条消息只会被其中一个消费者处理。此时应该使用TopicSubscriptionNameWithSuffix预置函数在 topic 名后追加后缀生成订阅名或者自定义一个GenerateSubscriptionName函数。例如上述示例中func(topic string) string { return test-sub_ topic }即为自定义命名策略的示范。反之如果你希望多个消费者各自独立地消费同一份完整消息流fanout 模式则需要为每个消费者使用不同的订阅名。这一点是 Google Cloud Pub/Sub 与 Kafka 消费组模型最大的差异之一消费组语义由订阅名体现而不是由消费者本身体现。Marshaler消息的序列化与反序列化Watermill 的message.Message无法被直接发送到 Google Cloud Pub/Sub——两者使用不同的消息格式因此需要经过marshaler的转换。在pkg/googlecloud包中你可以实现自己的 marshaler 来定制转换逻辑也可以直接使用默认实现DefaultMarshalerUnmarshaler官方文档内嵌了Marshaler接口与DefaultMarshalerUnmarshaler的定义源码。从默认实现的存在可以推断它负责完成*message.Message与 Google Cloud Pub/Sub 原生pubsub.Message及其 Attributes 元数据之间的双向转换从而在Publish和Subscribe两个方向上都保持一致的消息语义。对于大多数场景直接使用默认实现即可只有当需要自定义消息属性映射或 payload 编码时才需要实现自己的 marshaler。完整可运行示例将发布、订阅与命名策略串起来_examples/pubsubs/googlecloud/main.go 提供了一个可直接运行的完整程序该文件同时是 Watermill 官方 Getting Started 教程的源码package main import ( context log time github.com/ThreeDotsLabs/watermill github.com/ThreeDotsLabs/watermill-googlecloud/v2/pkg/googlecloud github.com/ThreeDotsLabs/watermill/message ) func main() { logger : watermill.NewStdLogger(false, false) subscriber, err : googlecloud.NewSubscriber( googlecloud.SubscriberConfig{ // custom function to generate Subscription Name, // there are also predefined TopicSubscriptionName and TopicSubscriptionNameWithSuffix available. GenerateSubscriptionName: func(topic string) string { return test-sub_ topic }, ProjectID: test-project, }, logger, ) if err ! nil { panic(err) } // Subscribe will create the subscription. Only messages that are sent after the subscription is created may be received. messages, err : subscriber.Subscribe(context.Background(), example.topic) if err ! nil { panic(err) } go process(messages) publisher, err : googlecloud.NewPublisher(googlecloud.PublisherConfig{ ProjectID: test-project, }, logger) if err ! nil { panic(err) } publishMessages(publisher) } 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)) // we need to Acknowledge that we received and processed the message, // otherwise, it will be resent over and over again. msg.Ack() } }运行方式本地模拟器使用仓库提供的 _examples/pubsubs/googlecloud/docker-compose.yml执行docker-compose up程序会通过PUBSUB_EMULATOR_HOST连接模拟器运行真实 GCP 环境设置GOOGLE_APPLICATION_CREDENTIALS环境变量指向服务账号凭据并将ProjectID改为真实项目 ID 后直接go run main.go。进阶实践从 Google Cloud Pub/Sub 到 MySQL 的事件管道Watermill 的 Google Cloud Pub/Sub 适配器不止用于简单的发布订阅演示也可以作为事件驱动架构中的事件摄取入口。仓库中的真实示例 _examples/real-world-examples/persistent-event-log/main.go 展示了如何将 Google Cloud Pub/Sub 与 Watermill Router 及 SQL 发布器组合构建一条Google Cloud Pub/Sub → MySQL的持久化事件日志管线subscriber : createSubscriber() // googlecloud.NewSubscriber(SubscriberConfig{ProjectID: example}, logger) publisher : createPublisher(db) // sql.NewPublisher(...) router.AddHandler( googlecloud-to-mysql, googleCloudTopic, // events subscriber, mysqlTable, // events publisher, func(msg *message.Message) ([]*message.Message, error) { consumedEvent : event{} err : json.Unmarshal(msg.Payload, consumedEvent) if err ! nil { return nil, err } log.Printf(received event %v with UUID %s, consumedEvent, msg.UUID) return []*message.Message{msg}, nil }, ) if err : router.Run(context.Background()); err ! nil { panic(err) }其中createSubscriber仅通过ProjectID字段创建订阅者此时订阅名采用默认的 topic 名事件由独立的simulateEvents协程以 JSON payload 形式持续发布到eventstopicRouter 中的处理器将其反序列化后转发到 MySQL 表。这个示例充分说明只要遵循 Watermill 的 Publisher/Subscriber 接口Google Cloud Pub/Sub 可以无缝地嵌入到 Router、中间件、插件构成的事件处理体系中与消息路由、重试、恢复等能力自由组合。小结Google Cloud Pub/Sub 适配器为 Watermill 带来了全托管、持久化、支持消费组ConsumerGroups的云端消息能力。回顾关键要点安装只需go get github.com/ThreeDotsLabs/watermill-googlecloud/v2连接完全由环境变量驱动生产环境用GOOGLE_APPLICATION_CREDENTIALS本地开发用PUBSUB_EMULATOR_HOST配合官方模拟器镜像无需安装 Cloud SDK订阅名即消费组语义多个消费者使用相同订阅名则协同消费每条消息只处理一次使用不同订阅名则各自消费完整消息流可通过TopicSubscriptionName、TopicSubscriptionNameWithSuffix或自定义GenerateSubscriptionName控制消息确认消费后必须调用msg.Ack()否则消息会被反复重投消息转换通过Marshaler默认实现为DefaultMarshalerUnmarshaler在 Watermill 消息与 Pub/Sub 原生消息之间双向转换能力边界支持持久化最长保留 7 天与消费组但不提供恰好一次投递与顺序保证架构设计时应据此选择合适的语义兜底方案。如需在本地快速体验直接使用仓库中的 _examples/pubsubs/googlecloud 示例并执行docker-compose up即可。【免费下载链接】watermillBuilding event-driven applications the easy way in Go.项目地址: https://gitcode.com/GitHub_Trending/wa/watermill创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门