GoWind 开源生态GoWind 开源生态
首页
框架
GoWind Admin
GoWind CMS
GoWind IM
GoWind UBA
GoWind IoT
GoWind Toolkit
GoWind Quant
GitHub
首页
框架
GoWind Admin
GoWind CMS
GoWind IM
GoWind UBA
GoWind IoT
GoWind Toolkit
GoWind Quant
GitHub
  • 介绍

    • GoWind 框架
    • 框架整体架构
  • go-wind 核心

    • go-wind 核心框架
    • App 生命周期管理
    • Context 传播
    • Transport 抽象
    • Log 门面接口
  • go-wind-plugins 插件

    • go-wind-plugins 插件总览
    • 插件配置系统
    • 插件注册机制(SPI)
    • 日志适配插件
    • 传输协议插件
    • 消息中间件插件
    • 编码解码插件
    • 安全与认证插件
    • 链路追踪插件
    • 缓存插件
    • 对象存储插件(OSS)
    • 限流插件
    • 指标监控插件
    • AI 插件
    • 工作流插件
    • 数据库与缓存插件
    • Bootstrap 集成与网关插件
    • 压缩与序列化插件
    • 模板渲染与验证插件
  • go-wind-bootstrap 启动器

    • go-wind-bootstrap 声明式启动器
    • Bootstrap 配置系统
    • Bootstrap SPI 机制
    • 声明式中间件编排
    • Bootstrap CLI 工具
    • Bootstrap 实战示例
  • 教程

    • 快速入门教程
    • 自定义插件开发教程
    • 多协议同时监听教程
    • 框架迁移指南

消息中间件插件

go-wind-plugins 提供多种消息队列(Broker)适配器,统一 broker.Broker 接口,支持发布/订阅模式。

一、Broker 接口

type Broker interface {
    Connect(ctx context.Context) error
    Disconnect(ctx context.Context) error
    Publish(ctx context.Context, topic string, msg *Message) error
    Subscribe(ctx context.Context, topic string, handler Handler) (Subscriber, error)
}

type Message struct {
    Headers map[string]string
    Body    []byte
}

二、适配器列表

适配器导入路径特点
Kafkaplugins/broker/kafka高吞吐、分区有序、消费者组
RabbitMQplugins/broker/rabbitmq路由灵活、ACK 确认
NATSplugins/broker/nats轻量、低延迟
Redisplugins/broker/redis简单、无额外依赖
Pulsarplugins/broker/pulsar多租户、持久化
RocketMQplugins/broker/rocketmq阿里系、事务消息
NSQplugins/broker/nsq去中心化、无 SPOF
MQTTplugins/broker/mqttIoT 协议、轻量

三、Kafka

import kafkaPlugin "github.com/tx7do/go-wind-plugins/broker/kafka"

broker := kafkaPlugin.NewBroker(
    kafkaPlugin.WithAddrs("localhost:9092"),
    kafkaPlugin.WithGroupID("my-service"),
    kafkaPlugin.WithVersion("3.0.0"),
)

发布消息

msg := &broker.Message{
    Headers: map[string]string{
        "trace_id": traceID,
    },
    Body: []byte(`{"event":"user_login","user_id":"123"}`),
}

broker.Publish(ctx, "user-events", msg)

订阅消息

sub, _ := broker.Subscribe(ctx, "user-events", func(ctx context.Context, msg *broker.Message) error {
    var event UserEvent
    json.Unmarshal(msg.Body, &event)
    processEvent(event)
    return nil  // 返回 nil 自动 ACK
})

defer sub.Unsubscribe()

YAML 配置

broker:
  kafka:
    addrs: ["localhost:9092"]
    group_id: "my-service"
    version: "3.0.0"
    topics:
      - name: user-events
        partitions: 6
        replication: 3
    consumer:
      initial_offset: latest    # latest | earliest
      session_timeout: 10s
      rebalance_timeout: 30s
    producer:
      acks: all                 # none | one | all
      compression: snappy       # none | gzip | snappy | lz4 | zstd
      batch_size: 16384

四、RabbitMQ

import rabbitmqPlugin "github.com/tx7do/go-wind-plugins/broker/rabbitmq"

broker := rabbitmqPlugin.NewBroker(
    rabbitmqPlugin.WithAddrs("amqp://guest:guest@localhost:5672/"),
    rabbitmqPlugin.WithExchange("events"),
    rabbitmqPlugin.WithExchangeType("topic"),
    rabbitmqPlugin.WithDurable(true),
    rabbitmqPlugin.WithQoS(10),       // prefetch count
)

YAML 配置

broker:
  rabbitmq:
    addrs: ["amqp://guest:guest@localhost:5672/"]
    exchange: events
    exchange_type: topic
    durable: true
    auto_delete: false
    qos:
      prefetch_count: 10
      prefetch_global: false
    queues:
      - name: order-events
        routing_key: order.*
        durable: true

五、NATS

import natsPlugin "github.com/tx7do/go-wind-plugins/broker/nats"

broker := natsPlugin.NewBroker(
    natsPlugin.WithAddrs("nats://localhost:4222"),
    natsPlugin.WithJetStream(true),
)

YAML 配置

broker:
  nats:
    addrs: ["nats://localhost:4222"]
    jetstream: true
    max_reconnects: 60
    reconnect_wait: 2s
    credentials_file: nats.creds

六、Redis Pub/Sub

import redisBrokerPlugin "github.com/tx7do/go-wind-plugins/broker/redis"

broker := redisBrokerPlugin.NewBroker(
    redisBrokerPlugin.WithAddr("localhost:6379"),
    redisBrokerPlugin.WithChannels("events", "notifications"),
)

YAML 配置

broker:
  redis:
    addr: localhost:6379
    db: 1
    channels:
      - events
      - notifications
    buffer_size: 1000

七、选择指南

需求推荐
日志/事件流Kafka(高吞吐、持久化)
任务队列RabbitMQ(ACK、路由)
微服务内部通信NATS(低延迟)
简单 Pub/SubRedis(无额外依赖)
IoT 设备消息MQTT(轻量、QoS)
金融/事务消息RocketMQ(事务消息)
大规模流处理Pulsar(多租户)

八、多 Broker 组合

broker:
  kafka:          # 主事件流
    addrs: ["localhost:9092"]
    group_id: "event-processor"
  rabbitmq:       # 任务队列
    addrs: ["amqp://localhost:5672/"]
    exchange: tasks
  redis:          # 通知广播
    addr: localhost:6379
    channels:
      - notifications
import (
    _ "github.com/tx7do/go-wind-plugins/broker/kafka"
    _ "github.com/tx7do/go-wind-plugins/broker/rabbitmq"
    _ "github.com/tx7do/go-wind-plugins/broker/redis"
)

相关文档

  • 插件配置系统
  • 插件总览
  • 编码解码插件
Edit this page
Last Updated:: 6/21/26, 9:28 PM
Contributors: Bobo
Prev
传输协议插件
Next
编码解码插件