首页
直播
壁纸
友链
搜索
1
MySQL如何解决深度分页问题?
11 阅读
2
buildadmin百度编辑器放两个只能显示一个
8 阅读
3
网站被 CC 攻击了?别慌,教你几招接地气的防护办法
7 阅读
4
thinkphp6 消息队列think-queue
6 阅读
5
microsoft store安装codex失败
5 阅读
服务器运维
后端技术
前端技术
梯子
数据库
小程序
登录
搜索
标签搜索
fastadmin
Redis
RabbitMQ
Go
服务器
codex
buildadmin
小程序
mysql
Nginx
Docker
Vue3
Node.js
MySQL优化
Linux
TypeScript
JWT
消息队列
Elasticsearch
搜索引擎
沿途的风景
累计撰写
37
篇文章
累计收到
0
条评论
首页
栏目
服务器运维
后端技术
前端技术
梯子
数据库
小程序
页面
直播
壁纸
友链
搜索到
2
篇与
» RabbitMQ
的结果
2026-07-24
RabbitMQ 消息队列入门与实战:从原理到 Spring Boot 整合
前言RabbitMQ 是基于 AMQP 协议的开源消息中间件,广泛应用于微服务架构中的异步解耦、削峰填谷和可靠投递场景。本文将从核心概念讲起,最终完成一个 Spring Boot 整合 RabbitMQ 的实战项目。一、核心概念1.1 AMQP 模型Producer → Exchange → (Binding) → Queue → Consumer概念说明Producer消息生产者Exchange交换机,接收消息并路由到队列Queue消息队列,存储消息Binding交换机与队列之间的绑定关系Routing Key路由键,Exchange 根据它决定消息去向Consumer消息消费者1.2 四种交换机类型类型说明适用场景Direct精确匹配 Routing Key点对点定向投递Fanout广播到所有绑定队列广播通知Topic模式匹配 Routing Key按主题订阅Headers根据消息头匹配复杂路由条件二、Docker 快速部署docker run -d --name rabbitmq \ -p 5672:5672 \ -p 15672:15672 \ -e RABBITMQ_DEFAULT_USER=admin \ -e RABBITMQ_DEFAULT_PASS=admin123 \ rabbitmq:3.13-management5672:AMQP 通信端口15672:管理界面端口,浏览器访问 http://localhost:15672三、Spring Boot 整合3.1 添加依赖<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency>3.2 配置spring: rabbitmq: host: localhost port: 5672 username: admin password: admin123 virtual-host: /3.3 队列与交换机配置@Configuration public class RabbitMQConfig { public static final String EXCHANGE = "order.exchange"; public static final String QUEUE = "order.queue"; public static final String ROUTING_KEY = "order.create"; @Bean public TopicExchange orderExchange() { return new TopicExchange(EXCHANGE); } @Bean public Queue orderQueue() { return QueueBuilder.durable(QUEUE) .withArgument("x-message-ttl", 60000) .withArgument("x-dead-letter-exchange", "order.dlx") .build(); } @Bean public Binding orderBinding() { return BindingBuilder.bind(orderQueue()) .to(orderExchange()) .with(ROUTING_KEY); } }3.4 生产者@Service public class OrderProducer { @Autowired private RabbitTemplate rabbitTemplate; public void sendOrder(OrderDTO order) { rabbitTemplate.convertAndSend( RabbitMQConfig.EXCHANGE, RabbitMQConfig.ROUTING_KEY, order, message -> { message.getMessageProperties() .setDeliveryMode(MessageDeliveryMode.PERSISTENT); return message; } ); } }3.5 消费者@Component public class OrderConsumer { @RabbitListener(queues = RabbitMQConfig.QUEUE) public void handleOrder(OrderDTO order, Channel channel, Message message) throws IOException { try { processOrder(order); channel.basicAck(message.getMessageProperties().getDeliveryTag(), false); } catch (Exception e) { channel.basicNack( message.getMessageProperties().getDeliveryTag(), false, true ); } } private void processOrder(OrderDTO order) { System.out.println("处理订单: " + order.getId()); } }3.6 配置手动 ACKspring: rabbitmq: listener: simple: acknowledge-mode: manual prefetch: 10 retry: enabled: true max-attempts: 3四、可靠性投递方案4.1 生产者确认(Publisher Confirm)@Configuration public class RabbitConfig { @Bean public RabbitTemplate rabbitTemplate(ConnectionFactory cf) { RabbitTemplate template = new RabbitTemplate(cf); template.setConfirmCallback((correlationData, ack, cause) -> { if (!ack) { log.error("消息未到达Exchange: {}", cause); } }); template.setReturnsCallback(returned -> { log.error("消息未路由到Queue: {}", returned.getMessage()); }); template.setMandatory(true); return template; } }4.2 死信队列(DLX)消息变成死信的条件:消息被消费者拒绝(basicNack + requeue=false)消息 TTL 过期队列达到最大长度@Bean public Queue deadLetterQueue() { return QueueBuilder.durable("order.dlq").build(); } @Bean public FanoutExchange deadLetterExchange() { return new FanoutExchange("order.dlx"); } @Bean public Binding dlxBinding() { return BindingBuilder.bind(deadLetterQueue()).to(deadLetterExchange()); }五、常见问题问题原因解决方案消息丢失消费者宕机未 ACK手动 ACK + 持久化消息重复网络抖动导致重复投递消费端幂等性(唯一ID去重)消息堆积消费速度小于生产速度增加消费者 / 限流队列阻塞单条消息处理过慢合理设置 prefetch总结RabbitMQ 在微服务架构中承担着异步解耦和流量削峰的关键角色。生产环境中务必做到:持久化 + 手动 ACK + 死信队列 + 消费幂等,才能保证消息的可靠投递。
2026年07月24日
0 阅读
0 评论
0 点赞
2026-07-24
Go + RabbitMQ 构建高并发任务队列实战
前言Go 语言的 goroutine 并发模型配合 RabbitMQ 的可靠消息投递,是构建高并发任务队列的经典方案。本文将实现一个完整的 Go + RabbitMQ 任务队列系统,包含生产者、消费者、重试机制和优雅退出。一、项目结构go-rabbitmq-queue/ ├── go.mod ├── config/ │ └── config.go ├── producer/ │ └── main.go ├── consumer/ │ └── main.go └── rabbitmq/ └── connection.go二、安装依赖go mod init github.com/yourname/go-rabbitmq-queue go get github.com/rabbitmq/amqp091-go三、RabbitMQ 连接封装package rabbitmq import ( "log" "time" amqp "github.com/rabbitmq/amqp091-go" ) type Config struct { URL string Exchange string Queue string RoutingKey string PrefetchCount int } type RabbitMQ struct { conn *amqp.Connection channel *amqp.Channel config Config } func New(cfg Config) (*RabbitMQ, error) { var conn *amqp.Connection var err error for i := 0; i < 5; i++ { conn, err = amqp.Dial(cfg.URL) if err == nil { break } log.Printf("连接失败(%d/5): %v", i+1, err) time.Sleep(3 * time.Second) } if err != nil { return nil, err } ch, err := conn.Channel() if err != nil { return nil, err } err = ch.ExchangeDeclare( cfg.Exchange, "direct", true, false, false, false, nil, ) if err != nil { return nil, err } args := amqp.Table{ "x-message-ttl": int32(60000), "x-dead-letter-exchange": cfg.Exchange + ".dlx", } _, err = ch.QueueDeclare( cfg.Queue, true, false, false, false, args, ) if err != nil { return nil, err } err = ch.QueueBind( cfg.Queue, cfg.RoutingKey, cfg.Exchange, false, nil, ) if err != nil { return nil, err } ch.Qos(cfg.PrefetchCount, 0, false) return &RabbitMQ{conn: conn, channel: ch, config: cfg}, nil } func (r *RabbitMQ) Channel() *amqp.Channel { return r.channel } func (r *RabbitMQ) Close() { r.channel.Close() r.conn.Close() }四、生产者package main import ( "encoding/json" "fmt" "log" "time" "github.com/yourname/go-rabbitmq-queue/rabbitmq" amqp "github.com/rabbitmq/amqp091-go" ) type Task struct { ID string `json:"id"` Type string `json:"type"` Payload interface{} `json:"payload"` } func main() { cfg := rabbitmq.Config{ URL: "amqp://admin:admin123@localhost:5672/", Exchange: "task.exchange", Queue: "task.queue", RoutingKey: "task.process", PrefetchCount: 10, } mq, err := rabbitmq.New(cfg) if err != nil { log.Fatal(err) } defer mq.Close() for i := 0; i < 100; i++ { task := Task{ ID: fmt.Sprintf("task-%d", i), Type: "email", Payload: map[string]string{ "to": fmt.Sprintf("user%d@example.com", i), "subject": "通知邮件", }, } body, _ := json.Marshal(task) err = mq.Channel().Publish( cfg.Exchange, cfg.RoutingKey, false, false, amqp.Publishing{ DeliveryMode: amqp.Persistent, ContentType: "application/json", Body: body, Timestamp: time.Now(), }, ) if err != nil { log.Printf("发送失败: %v", err) continue } log.Printf("发送任务: %s", task.ID) } log.Println("所有任务发送完成") }五、消费者(多 Worker 并发)package main import ( "encoding/json" "fmt" "log" "os" "os/signal" "sync" "syscall" "time" "github.com/yourname/go-rabbitmq-queue/rabbitmq" amqp "github.com/rabbitmq/amqp091-go" ) type Task struct { ID string `json:"id"` Type string `json:"type"` Payload interface{} `json:"payload"` } func processTask(task Task) error { log.Printf("处理任务: %s, 类型: %s", task.ID, task.Type) time.Sleep(500 * time.Millisecond) if time.Now().Unix()%10 == 0 { return fmt.Errorf("模拟处理失败") } log.Printf("任务完成: %s", task.ID) return nil } func startWorker(id int, mq *rabbitmq.RabbitMQ, wg *sync.WaitGroup) { defer wg.Done() msgs, err := mq.Channel().Consume( "task.queue", fmt.Sprintf("worker-%d", id), false, false, false, false, nil, ) if err != nil { log.Printf("Worker %d 启动失败: %v", id, err) return } log.Printf("Worker %d 启动", id) for msg := range msgs { var task Task if err := json.Unmarshal(msg.Body, &task); err != nil { log.Printf("Worker %d 解析失败: %v", id, err) msg.Nack(false, false) continue } if err := processTask(task); err != nil { log.Printf("Worker %d 处理失败: %s -> %v", id, task.ID, err) retryCount := getRetryCount(msg) if retryCount < 3 { msg.Nack(false, true) } else { log.Printf("Worker %d 任务 %s 重试超限,进入死信", id, task.ID) msg.Nack(false, false) } continue } msg.Ack(false) } log.Printf("Worker %d 退出", id) } func getRetryCount(msg amqp.Delivery) int { if deaths, ok := msg.Headers["x-death"].([]interface{}); ok && len(deaths) > 0 { if death, ok := deaths[0].(amqp.Table); ok { if count, ok := death["count"].(int64); ok { return int(count) } } } return 0 } func main() { cfg := rabbitmq.Config{ URL: "amqp://admin:admin123@localhost:5672/", Exchange: "task.exchange", Queue: "task.queue", RoutingKey: "task.process", PrefetchCount: 5, } mq, err := rabbitmq.New(cfg) if err != nil { log.Fatal(err) } defer mq.Close() var wg sync.WaitGroup workerCount := 5 for i := 1; i <= workerCount; i++ { wg.Add(1) go startWorker(i, mq, &wg) } sigs := make(chan os.Signal, 1) signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM) <-sigs log.Println("收到退出信号,等待 worker 完成...") mq.Close() wg.Wait() log.Println("所有 worker 已退出") }六、死信队列消费者func startDeadLetterConsumer(mq *rabbitmq.RabbitMQ) { ch := mq.Channel() ch.ExchangeDeclare("task.exchange.dlx", "fanout", true, false, false, false, nil) _, _ = ch.QueueDeclare("task.dlq", true, false, false, false, nil) _ = ch.QueueBind("task.dlq", "", "task.exchange.dlx", false, nil) msgs, _ := ch.Consume("task.dlq", "dlq-consumer", false, false, false, false, nil) go func() { for msg := range msgs { log.Printf("死信消息: %s", string(msg.Body)) msg.Ack(false) } }() }七、架构总结Producer → [task.exchange] → [task.queue] → 5个 Worker 并发消费 ↓ (失败/Nack) [task.exchange.dlx] → [task.dlq] → 死信消费者记录总结Go + RabbitMQ 的组合充分发挥了各自优势:Go 的 goroutine 让消费者可以轻松开几十个并发 worker,RabbitMQ 的 ACK 机制保证消息不丢失。生产环境注意:设置合理的 prefetch 避免消息堆积在单个 worker、实现死信队列处理失败消息、添加优雅退出逻辑避免消息中断。
2026年07月24日
0 阅读
0 评论
0 点赞
0:00