Skip to content

消息队列与异步通信

本篇是 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 级队列内成熟业务消息、复杂路由
Pulsarms 级分区内成长中大数据 + 业务混合
Redis Streamsms 级队列内简单轻量场景

本篇选 NATS 和 Kafka:NATS 简单易用,适合入门和云原生场景;Kafka 高吞吐,适合大数据场景。

二、异步通信模式

1. 队列(Queue / Competing Consumers)

多个消费者从同一队列消费,每条消息只被一个消费者处理。用于任务分发:

Producer ──► [Queue] ──► Consumer A
                    ──► Consumer B
                    ──► Consumer C

2. 发布订阅(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 -js
  • 4222:客户端连接端口
  • 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:PLAINTEXT
bash
docker-compose up -d

3. 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 的 RequiredAcksNone(不等)、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-workernotify-worker),处理逻辑不同。

启动后:

  1. POST http://localhost:8080/orders 创建订单。
  2. 三个 worker 各自收到消息并处理。
  3. 它们是不同的 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 流水线。

延伸阅读