GoFr 发布订阅(Pub/Sub)实战指南:多消息中间件接入、订阅发布与全链路分布式追踪
【免费下载链接】gofrAn opinionated GoLang framework for accelerated microservice development. Built in support for databases and observability.项目地址: https://gitcode.com/GitHub_Trending/go/gofr
发布订阅(Publisher-Subscriber)是 GoFr 内置支持的一种异步通信架构模式,它让消息的生产者与消费者彼此解耦,从而使每个组件都能按自身需求独立伸缩与维护。本文将以 GoFr 官方高级指南 docs/advanced-guide/using-publisher-subscriber/page.md 为主体,结合仓库源码与真实示例(examples/using-publisher、examples/using-subscriber),系统讲解 GoFr 中 Kafka、Google Pub/Sub、MQTT、NATS JetStream、Redis Pub/Sub、Azure Event Hubs、Amazon SQS 七类消息后端的配置接入方法,以及app.Subscribe订阅与ctx.GetPublisher().Publish发布两种核心 API 的用法,最后深入剖析 GoFr 为发布/订阅自动注入的分布式追踪机制。读完本文,你将能够在一个 GoFr 应用中零手工样板代码地完成异步事件驱动的生产与消费,并直接观测到跨服务、跨消息后端的完整调用链路。
认识发布订阅模式与 GoFr 的设计选择
发布订阅是一种用于不同实体之间异步通信的架构设计模式,这里的实体既可以是不同的应用,也可以是同一应用的多个实例。消息在组件之间传递时,组件彼此并不知晓对方的存在,也就是组件是解耦的。这种解耦让系统更灵活、更易伸缩,因为每个组件都可以按照自身需求独立扩容和维护。
在 GoFr 应用中,如果用户需要使用发布订阅模式,框架提供了多个消息中间件后端,包括Apache Kafka、Google PubSub、MQTT、NATS JetStream、Redis Pub/Sub、Azure Event Hubs 和 Amazon SQS。PubSub 的初始化发生在IoC 容器中,由容器统一处理 PubSub 客户端依赖,这样控制权就掌握在框架手中,从而提升了模块化、可测试性和可复用性。用户只需提供 topic 名称,即可在单个应用中对多个 topic 进行发布和订阅;通过容器暴露的GetPublisher与GetSubscriber方法即可拿到 Publisher 与 Subscriber 接口,完成读取单条消息或向消息中间件发布消息的操作。
从源码可以看到这两类接口的定义非常精简(pkg/gofr/datasource/pubsub/interface.go):
type Publisher interface { Publish(ctx context.Context, topic string, message []byte) error } type Subscriber interface { Subscribe(ctx context.Context, topic string) (*Message, error) }此外,完整的Client接口还组合了Publisher、Subscriber、健康检查、topic 管理(CreateTopic/DeleteTopic)、Query与Close(pkg/gofr/datasource/pubsub/interface.go),而Committer接口则负责消息提交(Commit()),它是消费语义(at-least-once 等)得以落实的关键。
Container 是 GoFr Context 的一部分。
在 pkg/gofr/container/container.go 中可以看到createPubSub会根据配置的PUBSUB_BACKEND分别创建 Kafka、Google、MQTT、Redis 等客户端,而GetPublisher/GetSubscriber直接返回容器内的PubSub客户端(pkg/gofr/container/container.go)。这意味着应用里只有一个共享的 PubSub 客户端,依赖注入由框架完成。
配置与初始化:PUBSUB_BACKEND 驱动后端选择
GoFr 通过环境变量/.env文件配置 PubSub 后端,其中最核心的变量是PUBSUB_BACKEND,它决定了应用将使用哪一个消息中间件。仓库自带的两个示例(examples/using-publisher/configs/.env、examples/using-subscriber/configs/.env)默认都使用 Kafka,并预置了切换到 MQTT、GOOGLE 的注释模板,可直接对照修改。
下面按中间件逐一说明配置与本地搭建方法。
Kafka
配置项
Kafka 是 GoFr 默认演示使用的后端,下表汇总了其全部配置项:
| 配置项 | 说明 | 必填 | 默认值 | 示例 | 合法格式 |
|---|---|---|---|---|---|
PUBSUB_BACKEND | 使用 Apache Kafka 作为消息中间件 | 是 | - | KAFKA | 非空字符串 |
PUBSUB_BROKER | Kafka broker 连接地址,多 broker 用逗号分隔 | 是 | - | localhost:9092或localhost:8087,localhost:8088,localhost:8089 | 非空字符串 |
CONSUMER_ID | 消费者组 ID,用于唯一标识消费者组 | 消费时需要 | - | order-consumer | 非空字符串 |
PUBSUB_OFFSET | 当消费者组在某个分区没有已提交偏移量时,从何处开始消费 | 否 | -1 | 10 | int |
KAFKA_BATCH_SIZE | 发送到分区前缓冲的消息条数上限 | 否 | 100 | 10 | 正整数 |
KAFKA_BATCH_BYTES | 发送到分区前单次请求的最大字节数 | 否 | 1048576 | 65536 | 正整数 |
KAFKA_BATCH_TIMEOUT | 未满批次消息刷入 Kafka 的时间上限(毫秒) | 否 | 1000 | 300 | 正整数 |
KAFKA_SECURITY_PROTOCOL | 与 Kafka 通信的安全协议(如 PLAINTEXT、SSL、SASL_PLAINTEXT、SASL_SSL) | 否 | PLAINTEXT | SASL_SSL | String |
KAFKA_SASL_MECHANISM | SASL 认证机制(如 PLAIN、SCRAM-SHA-256、SCRAM-SHA-512) | 否 | "" | PLAIN | String |
KAFKA_SASL_USERNAME | SASL 认证用户名 | 否 | "" | user | String |
KAFKA_SASL_PASSWORD | SASL 认证密码 | 否 | "" | password | String |
KAFKA_TLS_CERT_FILE | TLS 证书文件路径 | 否 | "" | /path/to/cert.pem | Path |
KAFKA_TLS_KEY_FILE | TLS 私钥文件路径 | 否 | "" | /path/to/key.pem | Path |
KAFKA_TLS_CA_CERT_FILE | TLS CA 证书文件路径 | 否 | "" | /path/to/ca.pem | Path |
KAFKA_TLS_INSECURE_SKIP_VERIFY | 是否跳过 TLS 证书校验 | 否 | false | true | Boolean |
一份完整的 Kafka.env配置示例如下:
PUBSUB_BACKEND=KAFKA# using apache kafka as message broker PUBSUB_BROKER=localhost:9092 CONSUMER_ID=order-consumer KAFKA_BATCH_SIZE=1000 KAFKA_BATCH_BYTES=1048576 KAFKA_BATCH_TIMEOUT=300 KAFKA_SASL_MECHANISM=PLAIN KAFKA_SASL_USERNAME=user KAFKA_SASL_PASSWORD=password KAFKA_TLS_CERT_FILE=/path/to/cert.pem KAFKA_TLS_KEY_FILE=/path/to/key.pem KAFKA_TLS_CA_CERT_FILE=/path/to/ca.pem KAFKA_TLS_INSECURE_SKIP_VERIFY=true其中PUBSUB_OFFSET默认值-1代表从最新的偏移量开始消费(示例 examples/using-subscriber/configs/.env 中则使用了-2,表示从最早偏移量开始,方便回放消息);KAFKA_BATCH_*三个参数直接控制生产端吞吐与延迟的权衡,批处理越大吞吐越高但单条消息延迟越大。这些配置在 pkg/gofr/container/container.go 的createKafkaPubSub中会被读取并组装进kafka.Config。
Docker 搭建 Kafka
使用 Bitnami 的 Kafka 3.4 镜像,以 KRaft 模式(无需 ZooKeeper)一键启动单节点:
docker run --name kafka-1 -p 9092:9092 \ -e KAFKA_ENABLE_KRAFT=yes \ -e KAFKA_CFG_PROCESS_ROLES=broker,controller \ -e KAFKA_CFG_CONTROLLER_LISTENER_NAMES=CONTROLLER \ -e KAFKA_CFG_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093 \ -e KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT \ -e KAFKA_CFG_ADVERTISED_LISTENERS=PLAINTEXT://127.0.0.1:9092 \ -e KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true \ -e KAFKA_BROKER_ID=1 \ -e KAFKA_CFG_CONTROLLER_QUORUM_VOTERS=1@127.0.0.1:9093 \ -e ALLOW_PLAINTEXT_LISTENER=yes \ -e KAFKA_CFG_NODE_ID=1 \ -v kafka_data:/bitnami \ bitnami/kafka:3.4注意KAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE=true开启了 topic 自动创建,便于快速验证;生产环境建议显式管理 topic。
Google Pub/Sub
配置项
PUBSUB_BACKEND=GOOGLE // using Google PubSub as message broker GOOGLE_PROJECT_ID=project-order // google projectId where the PubSub is configured GOOGLE_SUBSCRIPTION_NAME=order-consumer // unique subscription name to identify the subscribing entityGOOGLE_PROJECT_ID是 Google Cloud 中配置了 Pub/Sub 的项目 ID;GOOGLE_SUBSCRIPTION_NAME是标识订阅实体的唯一订阅名。此外还需要配置GOOGLE_APPLICATION_CREDENTIAL(指向服务账号 JSON 凭证),GoFr 底层通过 Google 官方客户端默认应用凭据(ADC)机制读取。
注意:在 Google Pub/Sub 中,一个订阅名只能访问一个 topic,框架会把 topic 名和订阅名拼接起来,在 Google 客户端上形成唯一的订阅名。
Docker 搭建 Google Pub/Sub 模拟器
docker pull gcr.io/google.com/cloudsdktool/google-cloud-cli:emulators docker run --name=gcloud-emulator -d -p 8086:8086 \ gcr.io/google.com/cloudsdktool/google-cloud-cli:emulators gcloud beta emulators pubsub start --project=test123 \ --host-port=0.0.0.0:8086启动模拟器后,应用连接时需设置PUBSUB_EMULATOR_HOST=localhost:8086指向本地模拟器,从而在无 Google 云账号的情况下完成开发与联调。
MQTT
配置项
PUBSUB_BACKEND=MQTT // using MQTT as pubsub MQTT_HOST=localhost // broker host URL MQTT_PORT=1883 // broker port MQTT_CLIENT_ID_SUFFIX=test // suffix to a random generated client-id(uuid v4) #some additional configs(optional) MQTT_PROTOCOL=tcp // protocol for connecting to broker can be tcp, tls, ws or wss MQTT_MESSAGE_ORDER=true // config to maintain/retain message publish order, by default this is false MQTT_USER=username // authentication username MQTT_PASSWORD=password // authentication passwordMQTT_HOST/MQTT_PORT:broker 的主机与端口;MQTT_CLIENT_ID_SUFFIX:框架会为客户端生成一个 uuid v4 作为 client id,该配置作为其后缀,用于多实例区分;MQTT_PROTOCOL:连接协议,支持tcp、tls、ws、wss;MQTT_MESSAGE_ORDER:是否保持/保留消息的发布顺序,默认false;MQTT_USER/MQTT_PASSWORD:认证用户名与密码。
注意:如果未提供
MQTT_HOST配置,应用会连接到 EMQX 提供的公共 MQTT 5 测试 broker。公共 broker 仅供本地开发验证,生产环境务必使用自有 broker 并开启认证。
Docker 搭建 MQTT(Eclipse Mosquitto)
docker run -d \ --name mqtt \ -p 8883:8883 \ -v \ eclipse-mosquitto:latest <path-to>/mosquitto.conf:/mosquitto/config/mosquitto.conf将本地编写的mosquitto.conf挂载进容器/mosquitto/config/mosquitto.conf,即可自定义端口、TLS 与认证规则。
NATS JetStream
NATS JetStream 是 GoFr 支持的外部 PubSub 提供方,也就是说,如果你不使用它,它不会被编译进你的二进制文件,从而保持二进制体积最小化。若需使用,可参考 NATS 官方文档了解 JetStream 的概念、连接与凭据(creds)配置。
配置项
PUBSUB_BACKEND=NATS PUBSUB_BROKER=nats://localhost:4222 NATS_STREAM=mystream NATS_SUBJECTS=orders.*,shipments.* NATS_MAX_WAIT=5s NATS_MAX_PULL_WAIT=500ms NATS_CONSUMER=my-consumer NATS_CREDS_FILE=/path/to/creds.json设置步骤
- 引入 NATS JetStream 外部驱动:
go get gofr.dev/pkg/gofr/datasource/pubsub/nats- 使用
AddPubSub方法将 NATS JetStream 驱动注入应用:
app := gofr.New() app.AddPubSub(nats.New(nats.Config{ Server: "nats://localhost:4222", Stream: nats.StreamConfig{ Stream: "mystream", Subjects: []string{"orders.*", "shipments.*"}, }, MaxWait: 5 * time.Second, MaxPullWait: 500 * time.Millisecond, Consumer: "my-consumer", CredsFile: "/path/to/creds.json", }))Docker 搭建 NATS
docker run -d \ --name nats \ -p 4222:4222 \ -p 8222:8222 \ -v \ nats:2.9.16 <path-to>/nats.conf:/nats/config/nats.conf4222是客户端连接端口,8222是 HTTP 监控管理端口。
配置项总览
| 名称 | 说明 | 必填 | 默认值 | 示例 |
|---|---|---|---|---|
PUBSUB_BACKEND | 设为NATS以使用 NATS JetStream 作为消息中间件 | 是 | - | NATS |
PUBSUB_BROKER | NATS 服务器 URL | 是 | - | nats://localhost:4222 |
NATS_STREAM | NATS stream 名称 | 是 | - | mystream |
NATS_SUBJECTS | 订阅的 subject 列表(逗号分隔) | 是 | - | orders.*,shipments.* |
NATS_MAX_WAIT | 批量请求的最大等待时间 | 否 | - | 5s |
NATS_MAX_PULL_WAIT | 单次 pull 请求的最大等待时间 | 否 | 0 | 500ms |
NATS_CONSUMER | NATS consumer 名称 | 否 | - | my-consumer |
NATS_CREDS_FILE | 认证凭据文件路径 | 否 | - | /path/to/creds.json |
使用提示:使用 NATS JetStream 订阅或发布时,请务必使用与你的 stream 配置相匹配的 subject 名称,否则消息将无法路由到对应 stream。
Redis Pub/Sub
Redis Pub/Sub 是一套轻量级的消息系统,GoFr 支持两种模式:
- Streams 模式(默认):基于 Redis Streams 实现持久化消息,支持消费者组与消息确认(acknowledgment),语义为at-least-once;
- PubSub 模式:标准 Redis Pub/Sub,即发即弃(fire-and-forget),无持久化,语义为at-most-once。
Redis 连接
Redis Pub/Sub 复用 Redis 数据源的连接配置(REDIS_HOST、REDIS_PORT、REDIS_DB、TLS 等),详见 docs/references/configs/page.md 中的 Redis 一节。
示例.env
PUBSUB_BACKEND=REDIS REDIS_HOST=localhost REDIS_PORT=6379 REDIS_USER=myuser REDIS_PASSWORD=mypassword REDIS_DB=0 REDIS_PUBSUB_DB=1 REDIS_TLS_ENABLED=true REDIS_TLS_CA_CERT=/path/to/ca.pem REDIS_TLS_CERT=/path/to/cert.pem REDIS_TLS_KEY=/path/to/key.pem # Streams mode (default) - requires consumer group REDIS_STREAMS_CONSUMER_GROUP=my-group REDIS_STREAMS_CONSUMER_NAME=my-consumer REDIS_STREAMS_BLOCK_TIMEOUT=5s REDIS_STREAMS_PEL_RATIO=0.7 # 70% PEL, 30% new messages REDIS_STREAMS_MAXLEN=1000 # To use PubSub mode instead, set: # REDIS_PUBSUB_MODE=pubsubDocker 搭建 Redis
无认证:
docker run -d \ --name redis \ -p 6379:6379 \ redis:7-alpine带密码认证:
docker run -d \ --name redis \ -p 6379:6379 \ redis:7-alpine redis-server --requirepass mypassword带 TLS:
docker run -d \ --name redis \ -p 6379:6379 \ -v /path/to/certs:/tls \ redis:7-alpine redis-server \ --tls-port 6380 \ --port 0 \ --tls-cert-file /tls/redis.crt \ --tls-key-file /tls/redis.key \ --tls-ca-cert-file /tls/ca.crtRedis Pub/Sub 专属配置
以下配置仅针对 Redis Pub/Sub 行为,基础连接/TLS 配置请参考 docs/references/configs/page.md 的 Redis 部分:
| 配置项 | 说明 | 默认值 | 示例 |
|---|---|---|---|
PUBSUB_BACKEND | 设为REDIS使用 Redis 作为 Pub/Sub 后端 | - | REDIS |
REDIS_PUBSUB_MODE | 模式:streams(默认,at-least-once)或pubsub(at-most-once) | streams | pubsub |
REDIS_STREAMS_CONSUMER_GROUP | 消费者组名(streams 模式必填) | - | mygroup |
REDIS_STREAMS_CONSUMER_NAME | 消费者名(可选,为空时自动生成) | - | consumer-1 |
REDIS_STREAMS_BLOCK_TIMEOUT | 流读取的阻塞超时。值越小(1s-2s)消息发现越快但 CPU 占用高;值越大(10s-30s)CPU 占用低但延迟高 | 5s | 2s或30s |
REDIS_STREAMS_PEL_RATIO | PEL(pending,待处理)消息与新增消息的读取比例(0.0-1.0)。该比例决定初始 PEL 分配,剩余容量始终用于填充新消息 | 0.7 | 0.5或0.8 |
REDIS_STREAMS_MAXLEN | 流的最大长度(近似裁剪)。设为0表示不限制 | 0(不限制) | 10000 |
REDIS_PUBSUB_DB | Pub/Sub 操作使用的 Redis 逻辑库。使用迁移(migrations)+ streams 模式时应与REDIS_DB区分开 | 15 | 1 |
REDIS_PUBSUB_BUFFER_SIZE | 消息缓冲区大小 | 100 | 1000 |
REDIS_PUBSUB_QUERY_TIMEOUT | Query 操作的超时时间 | 5s | 30s |
REDIS_PUBSUB_QUERY_LIMIT | Query 操作的消息条数上限 | 10 | 50 |
重要:如果
REDIS_STREAMS_CONSUMER_GROUP为空或未提供,尝试订阅时会直接报错;但发布操作不受影响,可以正常工作。
注意:topic 在首次发布时自动创建。当使用 GoFr 迁移功能配合 Streams 模式时,请保持
REDIS_DB与REDIS_PUBSUB_DB分离(默认分别为 0 和 15)。对于REDIS_STREAMS_BLOCK_TIMEOUT:实时场景建议 1s-2s,批量处理场景建议 10s-30s。
从源码 pkg/gofr/container/container.go 可以看到,Redis PubSub 通过redis.NewPubSub(conf, ...)构造函数初始化,且当 Redis PubSub(streams 模式)与主 Redis 数据源共用同一逻辑库时会输出警告(pkg/gofr/container/container.go),这正是文档强调“两个 DB 必须分开”的底层原因。
Azure Event Hubs
GoFr 从v1.22.0起支持 Azure Event Hubs。订阅时,GoFr 会读取配置中消费者组的全部分区,省去手动管理分区的麻烦。
设置步骤
Azure Event Hubs 同样是外部 PubSub 提供方——不使用它就不会被编入二进制。先引入外部驱动:
go get gofr.dev/pkg/gofr/datasource/pubsub/eventhub再用AddPubSub方法连接:
app := gofr.New() app.AddPubSub(eventhub.New(eventhub.Config{ ConnectionString: "Endpoint=sb://gofr-dev.servicebus.windows.net/;SharedAccessKeyName=RootManageSharedAccessKey;SharedAccessKey=<key>", ContainerConnectionString: "DefaultEndpointsProtocol=https;AccountName=gofrdev;AccountKey=<key>;EndpointSuffix=core.windows.net", StorageServiceURL: "https://gofrdev.windows.net/", StorageContainerName: "test", EventhubName: "test1", ConsumerGroup: "$Default", }))从 Event Hubs 订阅/发布时,请保持 topic 名称与 event-hub 名称一致。
配置说明
- 创建 Event Hubs 命名空间与 event hub;
- 由于 GoFr 需要记录“已读到哪、还剩哪些”的分区消费进度,因此需要依赖 Azure Blob 容器来持久化检查点(checkpoint),请提前创建好容器。
各必填字段的含义如下:
| 配置字段 | 含义 |
|---|---|
ConnectionString | Event Hubs 命名空间的连接字符串(主键) |
ContainerConnectionString | 用于存储检查点的 Azure 存储账户连接字符串 |
StorageServiceURL | 存储账户的 Blob Service URL |
StorageContainerName | 存储检查点的容器名称 |
EventhubName | 需要订阅/发布的 Event Hub 名称 |
Amazon SQS
GoFr 支持 Amazon Simple Queue Service(SQS)作为外部 PubSub 提供方。SQS 是托管的队列服务,用于解耦与伸缩微服务、分布式系统和 serverless 应用。
设置步骤
引入外部驱动:
go get gofr.dev/pkg/gofr/datasource/pubsub/sqs使用AddPubSub连接:
package main import ( "gofr.dev/pkg/gofr" "gofr.dev/pkg/gofr/datasource/pubsub/sqs" ) func main() { app := gofr.New() app.AddPubSub(sqs.New(&sqs.Config{ Region: "us-east-1", AccessKeyID: "your-access-key-id", // optional if using IAM roles SecretAccessKey: "your-secret-access-key", // optional if using IAM roles // Endpoint: "http://localhost:4566", // optional: for LocalStack })) app.Run() }注意:使用 IAM 角色(例如运行在 EC2 或 ECS 上)时,可以省略
AccessKeyID和SecretAccessKey,SDK 会自动使用实例的 IAM 角色凭据。
SQS 配置项
| 名称 | 说明 | 必填 | 默认值 | 示例 |
|---|---|---|---|---|
Region | SQS 队列所在的 AWS 区域 | 是 | - | us-east-1 |
AccessKeyID | AWS 访问密钥 ID | 否 | 使用默认凭据链 | AKIAIOSFODNN7EXAMPLE |
SecretAccessKey | AWS 秘密访问密钥 | 否 | 使用默认凭据链 | wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY |
SessionToken | 临时凭据的 AWS session token | 否 | - | FwoGZXIvYXdzE... |
Endpoint | 自定义 SQS 端点 URL,便于配合 LocalStack 本地开发 | 否 | AWS 默认端点 | http://localhost:4566 |
注意:SQS 队列必须在发布或订阅之前创建,可以使用 AWS CLI、AWS 控制台,或在迁移中通过
CreateTopic方法以编程方式创建队列。GoFr 默认支持标准队列,暂不支持 FIFO 队列;死信队列(DLQ)、广播(SNS)等高级特性可在基础设施层面配置。
LocalStack 本地开发
LocalStack 可以在本地模拟 AWS 服务,无需 AWS 账号即可开发与测试:
docker run -d \ --name localstack \ -p 4566:4566 \ -e SERVICES=sqs \ localstack/localstack:latestLocalStack 启动后,用 AWS CLI 创建队列:
aws --endpoint-url=http://localhost:4566 --region us-east-1 \ sqs create-queue --queue-name order-logs aws --endpoint-url=http://localhost:4566 --region us-east-1 \ sqs create-queue --queue-name products使用 LocalStack 时,将sqs.Config的Endpoint指向 LocalStack,并使用虚拟凭据:
app.AddPubSub(sqs.New(&sqs.Config{ Region: "us-east-1", Endpoint: "http://localhost:4566", AccessKeyID: "test", SecretAccessKey: "test", }))订阅(Subscribing)
在 GoFr 中添加订阅者与添加 HTTP handler 十分相似,这让开发可伸缩应用变得容易,因为订阅者与发送方/发布方完全解耦。用户可以定义一个订阅处理器完成消息处理,再通过app.Subscribe将其注入应用——这是典型的**控制反转(IoC)**模式,控制权留在框架手中,简化了开发与调试流程。
订阅处理器的签名如下:
func (ctx *gofr.Context) errorSubscribe方法会持续从配置的PUBSUB_BACKEND读取消息,后端可以是KAFKA、GOOGLE、MQTT、NATS、REDIS或AZURE_EVENTHUB;对于 NATS JetStream、Azure Event Hubs、Amazon SQS 等外部提供方,则使用app.AddPubSub()注入。这些配置都放在 configs 目录下的.env中。
处理器返回的 error 决定了哪些消息会被提交(commit)、哪些会被重新消费。
// First argument is the `topic name` followed by a handler which would process the // published messages continuously and asynchronously. app.Subscribe("order-status", func(ctx *gofr.Context)error{ // Handle the pub-sub message here })ctx为订阅者提供如下方法:
Bind():将消息值绑定到指定数据类型。消息可以转换为struct、map[string]any、int、bool、float64和string类型;Param(p string)/PathParam(p string):传入topic时返回当前 topic 名称。
从源码 pkg/gofr/datasource/pubsub/message.go 可以看到Bind的实现:它要求入参必须是指针(否则返回errNotPointer),并按目标类型分别走字符串、浮点、整数、布尔或 JSON 反序列化路径(bindStruct内部调用json.Unmarshal),这就是为什么Bind能同时支持基础类型与结构体。Param("topic")则直接返回消息的Topic字段(pkg/gofr/datasource/pubsub/message.go)。
订阅示例
package main import ( "gofr.dev/pkg/gofr" ) func main() { app := gofr.New() app.Subscribe("order-status", func(c *gofr.Context) error { var orderStatus struct { OrderId string `json:"orderId"` Status string `json:"status"` } err := c.Bind(&orderStatus) if err != nil { c.Logger.Error(err) // returning nil here as we would like to ignore the // incompatible message and continue reading forward return nil } c.Logger.Info("Received order ", orderStatus) return nil }) app.Run() }注意:当消息无法反序列化时,示例返回nil表示“忽略这条不兼容消息并继续向后读”;反之,如果你返回非nilerror,则该消息不会被提交,后续会被再次消费——这正是“返回错误决定重试语义”的落地方式。
仓库中的完整示例 examples/using-subscriber/main.go 在同一应用中订阅了两个 topic(products与order-logs),证明了“单应用多 topic 订阅”的能力。框架层的实现位于 pkg/gofr/subscriber.go:SubscriptionManager.startSubscriber会不断循环调用GetSubscriber().Subscribe,将读到的消息包装成gofr.Context交给 handler;handler 返回nil时自动调用msg.Commit()提交消息,出错时则延迟 2 秒后重试,且对 handler 内部的 panic 做了捕获恢复(panicRecovery),避免单个坏消息拖垮整个消费循环。
发布(Publishing)
发布操作建议在消息产生的位置完成。为此,用户可以从gofr.Context(ctx)中拿到发布接口来发送消息:
ctx.GetPublisher().Publish(ctx, "topic", msg)用户只需提供要发布的 topic 名称。GoFr 还支持向多个 topic 发布消息,因为应用往往需要向多个 topic 发送多种类型的消息。
发布示例
package main import ( "encoding/json" "gofr.dev/pkg/gofr" ) func main() { app := gofr.New() app.POST("/publish-order", order) app.Run() } func order(ctx *gofr.Context) (any, error) { type orderStatus struct { OrderId string `json:"orderId"` Status string `json:"status"` } var data orderStatus err := ctx.Bind(&data) if err != nil { return nil, err } msg, _ := json.Marshal(data) err = ctx.GetPublisher().Publish(ctx, "order-logs", msg) if err != nil { return nil, err } return "Published", nil }发布接口接收的是[]byte消息体(见 pkg/gofr/datasource/pubsub/interface.go),因此发布前通常先用json.Marshal序列化业务对象。仓库示例 examples/using-publisher/main.go 进一步展示了同一应用向order-logs、products两个 topic 分别发布消息的写法(POST /publish-order与POST /publish-product),并配合 migrations 在启动时自动创建 topic。
想直接跑通发布/订阅全流程,可对照 examples/using-publisher 与 examples/using-subscriber 两个示例:启动 Kafka 或 MQTT 后,在 subscriber 中注册对应 topic 的处理器,再用 curl 调用 publisher 的 HTTP 接口即可看到消息被消费的日志。
分布式追踪:发布/订阅全链路自动串联
GoFr 会自动对 Kafka、NATS JetStream、Google Pub/Sub、Amazon SQS 和 MQTT 的每一次发布与订阅调用生成追踪(trace),无需编写任何用户代码——只要配置了TRACE_EXPORTER(参见 docs/quick-start/observability/page.md 的 Tracing 一节),框架会自动完成一切埋点与上下文透传。
工作原理
当你调用ctx.GetPublisher().Publish(ctx, topic, msg)时,GoFr 会:
- 启动一个名为
<backend>-publish的 span(例如kafka-publish),SpanKind=Producer,并携带messaging.system、messaging.destination.name、messaging.operation=publish等语义属性; - 使用W3C Trace Context传播器把当前 trace 上下文注入到外发消息中:Kafka 中放在消息 headers 里,Google Pub/Sub 与 SQS 放在消息 attributes 里,NATS 放在消息 headers 中,以此类推。
当消息被app.Subscribe(topic, handler)注册的订阅者接收时,GoFr 会:
- 从入站消息中提取生产者的 trace 上下文;
- 启动名为
<backend>-subscribe的 span,SpanKind=Consumer,并且是生产者 span 的子 span——这意味着消费者 span 与发布者共享同一个TraceID,且父 span 指向发布者的 span; - 同时为生产者 span 附加一个 OpenTelemetryspan link。该 link 保留了扇出(fan-out)语义——同一条消息可能被多个消费者组消费——供 Jaeger、Tempo 等链路工具展示。
最终效果是:HTTP → publish → subscribe → publish → subscribe这样端到端的流程,在任意追踪 UI 中都会呈现为一条连通的 trace,瀑布图清晰可见:
[api-gateway ] POST /order (root) [api-gateway ] kafka-publish child of POST /order [order-service ] kafka-subscribe child of api-gateway's publish [+1 link] [order-service ] kafka-publish child of order-service's subscribe [notification-service ] kafka-subscribe child of order-service's publish [+1 link]MQTT 的追踪局限
MQTT 同样会被追踪,但它无法把发布端与订阅端接入同一条 trace,这一点在使用前值得了解,否则可能会误以为 exporter 出了问题。
原因在于:上述后端在线上消息中都留有承载 W3Ctraceparent的位置——Kafka 和 NATS 的 headers、SQS 与 Google Pub/Sub 的 message attributes。而 MQTT 3.1.1 的 PUBLISH 报文只有 topic、QoS、retain 标志和 payload,GoFr 使用的paho.mqtt.golang客户端也没有实现 MQTT 5 的 user properties——线上根本没有地方放置 trace 上下文。若强行放进 payload,则会破坏 topic 上所有非 GoFr 订阅者的消息,因此 GoFr 选择不这么做。
所以mqtt-publish与mqtt-subscribe与其它后端一样携带SpanKind、messaging.system、messaging.destination.name、messaging.operation语义属性,各自归属于调用它的请求的 trace——但消费者 span 不链接任何生产者,两者是分离的两条 trace。要跨 MQTT broker 跟踪消息,请在应用层用自己的消息标识进行关联。
采样与规模化
GoFr 的 tracer 使用ParentBased(TraceIDRatioBased(TRACER_RATIO))采样策略(见 pkg/gofr/otel.go)。由于消费者 span 继承生产者的采样决策,基于TRACER_RATIO的头部采样(head-based sampling)在整个链路上是一致的——如果生产者被采样丢弃,所有下游消费者 span 也会在创建时一并被丢弃。
对于高吞吐管道,建议将TRACER_RATIO设为小于1.0的值(例如0.1表示 10% 采样率)来控制 trace 体积;对于生命周期很长的异步 saga(一条 trace 保持数小时打开),则建议在 OpenTelemetry Collector 层面改用尾部采样(tail-based sampling)。
[!NOTE] 发布/订阅的分布式追踪是完全透明的——examples/using-publisher 与 examples/using-subscriber 中现有的示例无需任何追踪相关代码,跑起来就能产出连通的 trace。
小结
GoFr 把发布订阅从“选型 + 接线”简化成了“配一个.env+ 写一个 handler”:
- 后端选择:
PUBSUB_BACKEND一键切换 Kafka、Google Pub/Sub、MQTT、Redis Pub/Sub;NATS JetStream、Azure Event Hubs、Amazon SQS 通过app.AddPubSub作为外部驱动注入,不使用时不会膨胀二进制; - 消费:
app.Subscribe(topic, handler)以 IoC 方式注册订阅者,ctx.Bind统一完成消息反序列化,handler 返回值精确控制提交/重试语义(见 pkg/gofr/subscriber.go); - 生产:
ctx.GetPublisher().Publish(ctx, topic, msg)在消息产生处发送,天然支持多 topic; - 可观测:除 MQTT 外的后端均可自动串联端到端 trace,采样由
TRACER_RATIO统一控制。
无论是订单事件、日志流水还是库存变更,这套机制都能让你以最少的样板代码,构建出可伸缩、可观测的异步事件驱动架构。
【免费下载链接】gofrAn opinionated GoLang framework for accelerated microservice development. Built in support for databases and observability.项目地址: https://gitcode.com/GitHub_Trending/go/gofr
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考