Appearance
消息队列与异步通信
本篇是 Go 微服务系列的第九篇。前面几篇讲的通信都是同步的(HTTP、gRPC),但微服务架构中很多场景适合用异步通信:发通知、日志收集、订单状态变更、削峰填谷。消息队列(Message Queue, MQ)是异步通信的核心基础设施。本篇将讲解 MQ 在微服务中的作用、用 NATS 和 Kafka 实战、异步通信模式、可靠性保证、消费者组负载均衡,并给出一个完整的订单服务异步处理示例。
一、消息队列在微服务中的作用
1. 异步通信的价值
同步通信(HTTP/gRPC)的问题:
- 强耦合:调用方必须等被调用方响应。
- 级联故障:被调用方慢,调用方跟着慢,雪崩。
- 扩展性差:流量突增时双方都要扩容。
异步通信(MQ)的好处:
- 解耦:生产者只管发消息,不关心谁消费、何时消费。
- 削峰:流量洪峰先入队列,消费者按自己节奏处理。
- 广播:一条消息可被多个消费者独立消费。
- 重试:消费失败可重试,不丢数据。
2. 典型应用场景
- 通知 / 邮件:用户注册后异步发邮件、短信。
- 日志收集:各服务把日志发到 MQ,集中采集。
- 领域事件:订单创建后发事件,库存、积分、推荐各订阅。
- 任务分发:把任务发给 worker 池处理。
- 数据同步:数据库变更发到 MQ,下游同步到搜索引擎、缓存。
- 削峰填谷:秒杀场景请求先入队,后端按节奏处理。
3. 主流 MQ 对比
| MQ | 吞吐 | 延迟 | 持久化 | 顺序性 | 生态 | 典型场景 |
|---|---|---|---|---|---|---|
| Kafka | 极高 | ms 级 | 强 | 分区内 | 极丰富 | 大数据、日志 |
| NATS | 高 | 极低 | 可选 | 弱 | 云原生 | 微服务消息、IoT |
| RabbitMQ | 中 | μs 级 | 强 | 队列内 | 成熟 | 业务消息、复杂路由 |
| Pulsar | 高 | ms 级 | 强 | 分区内 | 成长中 | 大数据 + 业务混合 |
| Redis Streams | 中 | ms 级 | 弱 | 队列内 | 简单 | 轻量场景 |
本篇选 NATS 和 Kafka:NATS 简单易用,适合入门和云原生场景;Kafka 高吞吐,适合大数据场景。
二、异步通信模式
1. 队列(Queue / Competing Consumers)
多个消费者从同一队列消费,每条消息只被一个消费者处理。用于任务分发:
Producer ──► [Queue] ──► Consumer A
──► Consumer B
──► Consumer C2. 发布订阅(Pub/Sub)
每条消息被所有订阅者各自消费一份。用于事件广播:
Producer ──► [Topic] ──► Subscriber A (库存)
──► Subscriber B (积分)
──► Subscriber C (推荐)3. 请求响应(Request-Reply)
发送方发出消息后等待回复(通过 reply-to 队列)。模拟同步语义但走异步通道。NATS 内建支持。
4. 模式选择
| 需求 | 模式 |
|---|---|
| 任务分发 | 队列 |
| 事件广播 | 发布订阅 |
| 跨语言 RPC(异步) | 请求响应 |
| 数据流水线 | 队列 + 分区 |
三、NATS 实战:发布订阅
1. NATS 简介
NATS 是 Synadia 开源的高性能消息系统,特点:
- 极简:核心协议只有几个命令,开箱即用。
- 极快:单机百万级消息/秒。
- At-Most-Once:核心 NATS 不持久化,重启丢失(JetStream 补齐持久化)。
- 多模式:Pub/Sub、Queue Groups、Request-Reply。
- 云原生:常用于微服务通信。
2. 启动 NATS
Docker 启动:
bash
docker run -d --name nats \
-p 4222:4222 -p 8222:8222 \
nats:2.10 -js4222:客户端连接端口8222:监控端口-js:启用 JetStream(持久化)
3. 基础发布订阅
Go 客户端:github.com/nats-io/nats.go
go
package main
import (
"context"
"fmt"
"log"
"time"
"github.com/nats-io/nats.go"
)
func main() {
nc, err := nats.Connect("nats://127.0.0.1:4222")
if err != nil {
log.Fatal(err)
}
defer nc.Close()
// 订阅
sub, err := nc.SubscribeSync("orders.created")
if err != nil {
log.Fatal(err)
}
defer sub.Unsubscribe()
// 发布
for i := 0; i < 5; i++ {
msg := fmt.Sprintf(`{"order_id":%d,"ts":%d}`, i+1, time.Now().Unix())
if err := nc.Publish("orders.created", []byte(msg)); err != nil {
log.Fatal(err)
}
}
nc.Flush()
// 接收
for i := 0; i < 5; i++ {
msg, err := sub.NextMsgWithContext(context.Background())
if err != nil {
log.Fatal(err)
}
fmt.Printf("recv: subject=%s data=%s\n", msg.Subject, string(msg.Data))
}
}4. 异步订阅
实际项目用回调订阅,避免阻塞:
go
package main
import (
"fmt"
"log"
"sync"
"time"
"github.com/nats-io/nats.go"
)
func main() {
nc, _ := nats.Connect("nats://127.0.0.1:4222")
defer nc.Close()
var wg sync.WaitGroup
wg.Add(3)
// 多个订阅者:发布订阅模式,每个都收到
nc.Subscribe("orders.created", func(m *nats.Msg) {
fmt.Printf("[库存] recv: %s\n", string(m.Data))
wg.Done()
})
nc.Subscribe("orders.created", func(m *nats.Msg) {
fmt.Printf("[积分] recv: %s\n", string(m.Data))
wg.Done()
})
nc.Subscribe("orders.created", func(m *nats.Msg) {
fmt.Printf("[推荐] recv: %s\n", string(m.Data))
wg.Done()
})
time.Sleep(100 * time.Millisecond) // 等订阅就绪
nc.Publish("orders.created", []byte(`{"order_id":1}`))
nc.Flush()
wg.Wait()
log.Println("done")
}5. Queue Group(竞争消费)
把订阅者放到一个 queue group 里,每条消息只被其中一个消费:
go
package main
import (
"fmt"
"log"
"sync"
"time"
"github.com/nats-io/nats.go"
)
func main() {
nc, _ := nats.Connect("nats://127.0.0.1:4222")
defer nc.Close()
var wg sync.WaitGroup
wg.Add(10)
// 3 个 worker 在同一 queue group
for i := 1; i <= 3; i++ {
workerID := i
nc.QueueSubscribe("tasks", "workers", func(m *nats.Msg) {
fmt.Printf("[worker-%d] processed: %s\n", workerID, string(m.Data))
wg.Done()
})
}
time.Sleep(100 * time.Millisecond)
for i := 0; i < 10; i++ {
nc.Publish("tasks", []byte(fmt.Sprintf("task-%d", i)))
}
nc.Flush()
wg.Wait()
log.Println("all tasks done")
}6. JetStream 持久化
核心 NATS 不持久化,重启即丢。JetStream 提供持久化、ACK、重放:
go
package main
import (
"fmt"
"log"
"time"
"github.com/nats-io/nats.go"
)
func main() {
nc, _ := nats.Connect("nats://127.0.0.1:4222")
defer nc.Close()
js, _ := nc.JetStream()
// 创建 stream(持久化的 subject 集合)
js.AddStream(&nats.StreamConfig{
Name: "ORDERS",
Subjects: []string{"orders.*"},
Storage: nats.FileStorage,
})
// 发布持久消息
for i := 0; i < 5; i++ {
_, err := js.Publish("orders.created", []byte(fmt.Sprintf("order-%d", i)))
if err != nil {
log.Fatal(err)
}
}
// 持久订阅
sub, err := js.Subscribe("orders.created", func(m *nats.Msg) {
fmt.Printf("recv: %s\n", string(m.Data))
m.Ack()
}, nats.Durable("order-processor"), nats.ManualAck())
if err != nil {
log.Fatal(err)
}
time.Sleep(time.Second)
sub.Unsubscribe()
}四、Kafka 实战:高吞吐消息处理
1. Kafka 简介
Kafka 是 LinkedIn 开源的分布式流平台,核心概念:
- Topic:消息主题,类似 NATS 的 subject。
- Partition:分区,一个 Topic 可分多个分区,并行消费。
- Producer:生产者。
- Consumer / Consumer Group:消费者组,组内分区独占消费。
- Broker:Kafka 服务节点。
- Offset:消费者在分区中的位置。
特点:高吞吐(单集群百万 QPS)、持久化、可重放、分区顺序。
2. 启动 Kafka
Docker Compose 启动单节点(带 KRaft,无需 ZooKeeper):
yaml
# docker-compose.yml
services:
kafka:
image: bitnami/kafka:3.7
ports:
- "9092:9092"
environment:
KAFKA_CFG_NODE_ID: 1
KAFKA_CFG_PROCESS_ROLES: controller,broker
KAFKA_CFG_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_CFG_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_CFG_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_CFG_CONTROLLER_QUORUM_VOTERS: 1@localhost:9093
KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXTbash
docker-compose up -d3. Go 客户端
Go 生态主流 Kafka 客户端:github.com/segmentio/kafka-go(轻量)和 github.com/IBM/sarama(功能全)。本篇用 kafka-go。
生产者:
go
package main
import (
"context"
"fmt"
"log"
"time"
"github.com/segmentio/kafka-go"
)
func main() {
writer := &kafka.Writer{
Addr: kafka.TCP("localhost:9092"),
Topic: "orders",
Balancer: &kafka.LeastBytes{},
BatchTimeout: 10 * time.Millisecond,
RequiredAcks: kafka.RequireAll,
}
defer writer.Close()
for i := 0; i < 10; i++ {
err := writer.WriteMessages(context.Background(), kafka.Message{
Key: []byte(fmt.Sprintf("order-%d", i)),
Value: []byte(fmt.Sprintf(`{"order_id":%d,"amount":%d}`, i, i*100)),
})
if err != nil {
log.Printf("write failed: %v", err)
continue
}
}
log.Println("all messages sent")
}消费者:
go
package main
import (
"context"
"fmt"
"log"
"github.com/segmentio/kafka-go"
)
func main() {
r := kafka.NewReader(kafka.ReaderConfig{
Brokers: []string{"localhost:9092"},
Topic: "orders",
GroupID: "order-processor",
MinBytes: 1,
MaxBytes: 10 * 1024 * 1024,
CommitInterval: 1 * 1000 * 1000 * 1000, // 1s 自动提交
})
defer r.Close()
ctx := context.Background()
for {
m, err := r.ReadMessage(ctx)
if err != nil {
log.Printf("read failed: %v", err)
break
}
fmt.Printf("partition=%d offset=%d key=%s value=%s\n",
m.Partition, m.Offset, string(m.Key), string(m.Value))
}
}五、消息可靠性保证
1. 三种交付语义
- At-Most-Once:消息最多交付一次,可能丢。核心 NATS 默认。
- At-Least-Once:至少一次,可能重复。Kafka 默认。
- Exactly-Once:恰好一次,最难实现。Kafka 事务支持。
2. 生产端可靠性
- acks:Kafka 的
RequiredAcks:None(不等)、Leader(leader 写入即可)、All(所有副本确认)。 - 重试:发送失败重试,注意幂等。
- 同步发送:关键场景用
WriteMessages同步等待 ack。
3. 消费端可靠性
- 手动提交 offset:处理成功后才提交,避免漏处理。
- 幂等消费:因为 At-Least-Once 会重复,消费逻辑必须幂等。
- 死信队列(DLQ):处理失败的消息进 DLQ,避免阻塞主流。
手动提交示例:
go
package main
import (
"context"
"log"
"github.com/segmentio/kafka-go"
)
func main() {
r := kafka.NewReader(kafka.ReaderConfig{
Brokers: []string{"localhost:9092"},
Topic: "orders",
GroupID: "order-processor",
MinBytes: 1,
MaxBytes: 10 * 1024 * 1024,
})
defer r.Close()
ctx := context.Background()
for {
m, err := r.FetchMessage(ctx)
if err != nil {
log.Printf("fetch: %v", err)
continue
}
// 处理消息
if err := processOrder(m.Value); err != nil {
log.Printf("process failed, sending to DLQ: %v", err)
// 实际项目:发送到死信 topic
}
// 处理成功后提交
if err := r.CommitMessages(ctx, m); err != nil {
log.Printf("commit: %v", err)
}
}
}
func processOrder(data []byte) error {
log.Printf("processing: %s", string(data))
return nil
}4. 幂等消费
常见做法:用业务唯一 ID 去重:
go
package main
import (
"context"
"encoding/json"
"log"
"sync"
)
type OrderEvent struct {
OrderID string `json:"order_id"`
Amount int `json:"amount"`
}
type IdempotentProcessor struct {
processed sync.Map
}
func (p *IdempotentProcessor) Process(ctx context.Context, data []byte) error {
var evt OrderEvent
if err := json.Unmarshal(data, &evt); err != nil {
return err
}
if _, loaded := p.processed.LoadOrStore(evt.OrderID, true); loaded {
log.Printf("order %s already processed, skip", evt.OrderID)
return nil
}
// 真实处理
log.Printf("processed order %s amount=%d", evt.OrderID, evt.Amount)
return nil
}生产环境用 Redis SET NX 代替 sync.Map,跨实例去重。
六、消费者组的负载均衡
1. Kafka 消费者组
Kafka 的消费者组机制:
- 同一 GroupID 的消费者分摊 Topic 的所有分区。
- 每个分区同一时刻只被组内一个消费者消费。
- 消费者数 > 分区数时,多出的消费者闲置。
- 消费者增减时触发 rebalance,分区重新分配。
Topic: orders (3 partitions)
│
Consumer Group: order-processor
├── Consumer A ← partition 0, 1
└── Consumer B ← partition 2扩容消费者时,分区会被重新分配,吞吐随之提升。
2. Rebalance 的代价
Rebalance 期间消费者暂停消费,可能造成秒级延迟。优化策略:
- 合理设置 session timeout 和 heartbeat。
- 使用 Cooperative Rebalance(Kafka 2.4+)减少重平衡范围。
- 避免频繁上下线。
3. NATS 的 Queue Group
NATS 的 Queue Group 是类似机制:同一 queue 名下的订阅者分摊消息。NATS 没有 partition 概念,但 JetStream 的 consumer 同样支持。
七、完整示例:订单服务的异步处理
下面给出一个完整的订单服务:HTTP 接收订单 → 发到 Kafka → 多个消费者处理(库存、积分、通知)。
order-api/main.go(接收订单并生产消息):
go
package main
import (
"context"
"encoding/json"
"log"
"net/http"
"time"
"github.com/google/uuid"
"github.com/segmentio/kafka-go"
)
type Order struct {
OrderID string `json:"order_id"`
UserID string `json:"user_id"`
ProductID string `json:"product_id"`
Amount float64 `json:"amount"`
}
type OrderService struct {
writer *kafka.Writer
}
func NewOrderService(kafkaAddr, topic string) *OrderService {
return &OrderService{
writer: &kafka.Writer{
Addr: kafka.TCP(kafkaAddr),
Topic: topic,
Balancer: &kafka.LeastBytes{},
BatchTimeout: 10 * time.Millisecond,
RequiredAcks: kafka.RequireAll,
},
}
}
func (s *OrderService) Create(ctx context.Context, o *Order) error {
o.OrderID = uuid.NewString()
data, _ := json.Marshal(o)
return s.writer.WriteMessages(ctx, kafka.Message{
Key: []byte(o.OrderID),
Value: data,
})
}
func (s *OrderService) Close() error { return s.writer.Close() }
func main() {
svc := NewOrderService("localhost:9092", "orders")
defer svc.Close()
mux := http.NewServeMux()
mux.HandleFunc("/orders", func(w http.ResponseWriter, r *http.Request) {
if r.Method != http.MethodPost {
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
return
}
var o Order
if err := json.NewDecoder(r.Body).Decode(&o); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
return
}
if err := svc.Create(r.Context(), &o); err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
w.WriteHeader(http.StatusAccepted)
_ = json.NewEncoder(w).Encode(map[string]string{"order_id": o.OrderID})
})
log.Println("order api on :8080")
log.Fatal(http.ListenAndServe(":8080", mux))
}inventory-worker/main.go(库存消费者):
go
package main
import (
"context"
"encoding/json"
"log"
"github.com/segmentio/kafka-go"
)
type Order struct {
OrderID string `json:"order_id"`
UserID string `json:"user_id"`
ProductID string `json:"product_id"`
Amount float64 `json:"amount"`
}
func main() {
r := kafka.NewReader(kafka.ReaderConfig{
Brokers: []string{"localhost:9092"},
Topic: "orders",
GroupID: "inventory-worker",
MinBytes: 1,
MaxBytes: 10 * 1024 * 1024,
})
defer r.Close()
log.Println("inventory worker started")
for {
m, err := r.ReadMessage(context.Background())
if err != nil {
log.Printf("read: %v", err)
continue
}
var o Order
if err := json.Unmarshal(m.Value, &o); err != nil {
log.Printf("unmarshal: %v", err)
continue
}
if err := deductStock(o.ProductID); err != nil {
log.Printf("deduct stock failed: %v", err)
continue
}
log.Printf("[库存] order %s 扣减库存成功", o.OrderID)
}
}
func deductStock(productID string) error {
// 实际项目调用库存服务或 DB
return nil
}points-worker/main.go(积分消费者)和 notify-worker/main.go(通知消费者)结构相同,只是 GroupID 不同(points-worker、notify-worker),处理逻辑不同。
启动后:
POST http://localhost:8080/orders创建订单。- 三个 worker 各自收到消息并处理。
- 它们是不同的 consumer group,互不影响,可独立扩容。
这种架构的好处:
- 解耦:订单服务只管发消息,不关心下游。
- 独立扩容:哪个 worker 慢就扩哪个。
- 故障隔离:积分服务挂了不影响库存扣减。
- 可扩展:新增订阅者无需修改订单服务。
八、消息顺序性
1. 顺序性的难度
分布式 MQ 要保证全局顺序极难(牺牲并行),通常退而求其次:分区内顺序。
2. Kafka 分区顺序
Kafka 同一分区内的消息严格有序(按写入顺序)。所以:
- 把同一业务实体的消息路由到同一分区(如按
order_id哈希)。 - 消费者按分区串行消费,保证顺序。
go
writer := &kafka.Writer{
Addr: kafka.TCP("localhost:9092"),
Topic: "orders",
Balancer: &kafka.Hash{}, // 按 key 哈希到分区
}
// key 用 order_id,保证同一订单的事件顺序
writer.WriteMessages(ctx, kafka.Message{
Key: []byte(orderID),
Value: data,
})3. NATS JetStream
NATS JetStream 也支持按 subject 分区(subjects: orders.*),同一 subject 的消息保序。
九、小结
本篇我们学习了消息队列与异步通信:
- 价值:解耦、削峰、广播、重试,应对同步通信的痛点。
- 三大模式:队列(竞争消费)、发布订阅(广播)、请求响应。
- NATS:极简、极快,核心 NATS 不持久化,JetStream 提供持久化和 ACK。
- Kafka:高吞吐、分区并行、消费者组负载均衡,适合大数据和业务混合场景。
- 可靠性:生产端 acks + 重试,消费端手动提交 + 幂等 + 死信队列。
- 消费者组:组内分区分摊,扩容时触发 rebalance,注意减少 rebalance 代价。
- 顺序性:分区内顺序,按业务 ID 哈希路由到同一分区。
- 完整示例:订单服务 HTTP → Kafka → 多个消费者独立处理,体现解耦和独立扩展。
下一篇我们将进入微服务部署与编排,学习 Docker 多阶段构建、docker-compose、Kubernetes、Helm、CI/CD 流水线。
延伸阅读:
- NATS 文档:https://docs.nats.io/
- Kafka 文档:https://kafka.apache.org/documentation/
- kafka-go:https://github.com/segmentio/kafka-go
- 《Designing Event-Driven Systems》Ben Stopford 著