【原创】分布式之消息队列复习精讲
【原创】分布式之消息队列复习精讲在分布式系统中消息队列Message Queue, MQ是解耦、异步、流量削峰的核心组件。无论是微服务通信、事件驱动架构还是大数据处理MQ 都扮演着关键角色。本文将从实战出发带你快速复习消息队列的核心概念、常见选型、以及如何用代码实现一个简易版 MQ。## 为什么需要消息队列假设你有一个电商系统用户下单后需要执行创建订单、扣减库存、发送通知、更新积分。如果所有操作同步执行高并发下数据库会瞬间崩溃。使用 MQ 后订单系统发送“下单成功”消息到队列其他服务异步消费系统吞吐量提升数倍。## 核心概念速览-生产者发送消息的一方。-消费者接收并处理消息的一方。-队列存储消息的缓冲区通常支持持久化。-主题Topic逻辑分类如“订单”、“支付”。-ACK机制消费者处理成功后通知队列删除消息防止丢消息。## 常见消息队列对比| 组件 | 特点 | 适用场景 ||-----------|--------------------------|------------------------|| RabbitMQ | 功能丰富支持多种协议 | 企业级应用复杂路由 || Kafka | 高吞吐持久化分布式 | 日志收集流处理 || Redis Stream | 轻量级内存型 | 简单异步任务 |## 实战用 Python 实现一个简易消息队列为了深入理解原理我们用 Python 实现一个基于内存 多线程的 MQ支持生产者和消费者模式。### 1. 简易内存队列pythonimport threadingimport timeimport queueclass SimpleMQ: 简易消息队列支持多生产者、多消费者基于线程安全队列 def __init__(self, maxsize100): self.queue queue.Queue(maxsizemaxsize) # 线程安全队列 self.consumers [] # 消费者列表 def produce(self, message): 生产者发送消息 try: self.queue.put(message, blockFalse) # 非阻塞写入 print(f[生产者] 发送: {message}) except queue.Full: print([生产者] 队列已满消息丢失) def consume(self, consumer_id): 消费者持续消费 while True: try: message self.queue.get(timeout1) # 阻塞1秒 print(f[消费者 {consumer_id}] 处理: {message}) # 模拟处理耗时 time.sleep(0.5) self.queue.task_done() # 通知队列任务完成 except queue.Empty: print(f[消费者 {consumer_id}] 等待消息...) time.sleep(0.1)# 使用示例if __name__ __main__: mq SimpleMQ(maxsize5) # 启动2个消费者线程 for i in range(2): t threading.Thread(targetmq.consume, args(i,)) t.daemon True # 设置为守护线程主线程结束即退出 t.start() # 生产者发送10条消息 for i in range(10): mq.produce(f订单_{i}) time.sleep(0.2) # 模拟生产间隔 time.sleep(3) # 等待消费者处理完 print(主线程结束)输出示例[生产者] 发送: 订单_0[消费者 0] 处理: 订单_0[消费者 1] 等待消息...[生产者] 发送: 订单_1[消费者 1] 处理: 订单_1...原理分析- 使用queue.Queue保证线程安全内部有锁机制。- 生产者和消费者通过队列解耦支持异步处理。- 实际生产环境需考虑持久化、分布式、ACK机制等。## 实战基于 Redis 的分布式消息队列Redis 的 Stream 数据结构天然支持消息队列适合轻量级场景。下面演示如何用 Redis Stream 实现分布式 MQ。### 2. Redis Stream 实现 MQpythonimport redisimport timeimport threadingclass RedisMQ: 基于 Redis Stream 的分布式消息队列 需要安装 redis-py: pip install redis def __init__(self, stream_nameorder_stream, group_nameorder_group): self.client redis.Redis(hostlocalhost, port6379, decode_responsesTrue) self.stream stream_name self.group group_name # 创建消费者组如果不存在 try: self.client.xgroup_create(self.stream, self.group, id0, mkstreamTrue) except redis.exceptions.ResponseError as e: if BUSYGROUP not in str(e): raise e def produce(self, message): 生产者发送消息到 Stream msg_id self.client.xadd(self.stream, message) print(f[生产者] 发送成功ID: {msg_id}) return msg_id def consume(self, consumer_name): 消费者从消费者组消费消息 while True: try: # 读取未确认消息block1000mscount1 results self.client.xreadgroup( self.group, consumer_name, {self.stream: }, count1, block1000 ) if results: for stream, messages in results: for msg_id, msg_data in messages: print(f[消费者 {consumer_name}] 处理: {msg_data}) # 模拟处理 time.sleep(0.5) # 确认消息已处理 self.client.xack(self.stream, self.group, msg_id) print(f[消费者 {consumer_name}] ACK: {msg_id}) else: time.sleep(0.1) except Exception as e: print(f[消费者 {consumer_name}] 错误: {e}) time.sleep(1)# 使用示例if __name__ __main__: mq RedisMQ() # 启动2个消费者不同消费者名 for i in range(2): t threading.Thread(targetmq.consume, args(fconsumer_{i},)) t.daemon True t.start() # 生产者发送消息 for i in range(5): mq.produce({order_id: forder_{i}, amount: i*100}) time.sleep(0.3) time.sleep(5) print(主线程结束)关键点-消费者组允许多个消费者分摊消息保证一条消息只被一个消费者处理。-XACK 机制消费者处理完后必须确认否则消息会重新投递防止丢失。-持久化Redis 将 Stream 数据持久化到磁盘重启不丢失。## 消息队列的常见问题与解决方案### 1. 消息丢失-生产者开启确认模式如 RabbitMQ 的publisher confirms。-队列持久化消息Redis 用 RDB/AOFKafka 用副本机制。-消费者手动 ACK处理完再确认。### 2. 重复消费- 保证幂等性消费时检查业务唯一键如订单号是否已处理。- 例如数据库插入时使用ON DUPLICATE KEY UPDATE。### 3. 消息积压- 增加消费者数量水平扩展。- 临时扩容用 Kafka 或 RabbitMQ 的动态扩缩容机制。## 总结消息队列是分布式系统的“胶水”解决了同步阻塞、服务耦合、流量冲击三大痛点。通过本文的简易代码实现你应该理解了1.核心模型生产者 → 队列 → 消费者本质是生产者-消费者模式的分布式化。2.技术选型根据吞吐量、可靠性、运维成本选择 RabbitMQ、Kafka 或 Redis Stream。3.实战要点注意 ACK 机制、幂等性、持久化配置避免消息丢失或重复。建议在生产环境中优先使用成熟的消息队列如 Kafka 用于大数据、RabbitMQ 用于传统业务并配合监控如 Prometheus Grafana观察队列深度和消费延迟。复习至此你已经掌握了消息队列的核心知识可以自信应对分布式系统设计面试

相关新闻

104、Arduino Nano 33 BLE Sense的语音识别案例

104、Arduino Nano 33 BLE Sense的语音识别案例

104、Arduino Nano 33 BLE Sense的语音识别案例 从一次深夜调试说起 凌晨两点,我盯着串口监视器里不断跳出的“Unknown”字样,差点把这块Nano 33 BLE Sense扔出窗外。麦克风焊盘被我反复按压了十几遍,焊锡都快磨亮了,结果还是识别不出“yes”和“no”这两个最简单的词。后…

2026/7/26 1:09:58 阅读更多 →
扣子文件处理机器人部署避坑清单(2024最新版):从权限配置到异常熔断全链路解析

扣子文件处理机器人部署避坑清单(2024最新版):从权限配置到异常熔断全链路解析

更多请点击: https://kaifayun.com 第一章:扣子文件处理机器人部署避坑清单(2024最新版):从权限配置到异常熔断全链路解析 权限配置的三大致命误区 扣子平台对文件类机器人强制要求细粒度权限校验。常见错误包括&…

2026/7/26 1:08:58 阅读更多 →
3步实战:用SyncTrayzor轻松搭建Windows跨设备文件同步系统

3步实战:用SyncTrayzor轻松搭建Windows跨设备文件同步系统

3步实战:用SyncTrayzor轻松搭建Windows跨设备文件同步系统 【免费下载链接】SyncTrayzor Windows tray utility / filesystem watcher / launcher for Syncthing 项目地址: https://gitcode.com/gh_mirrors/sy/SyncTrayzor SyncTrayzor是Windows系统上最优雅…

2026/7/26 1:07:57 阅读更多 →

最新新闻

缠论画线N 同花顺期货通指标

缠论画线N 同花顺期货通指标

今天给大家带来是一款同花顺期货通指标,并且已经上架到同花顺期货通的指标广场上了。喜欢的朋友可以去指标广场安装试用!!友情提示:(指标只是辅助,不作建议)拼多多店铺:指标公式编写…

2026/7/26 1:18:01 阅读更多 →
嵌入式音频开发实战:深入解析McASP寄存器配置与I2S数据流管理

嵌入式音频开发实战:深入解析McASP寄存器配置与I2S数据流管理

1. 项目概述与I2S/McASP核心价值在嵌入式音频系统开发中,无论是智能音箱、车载娱乐系统还是专业音频设备,数字音频数据的可靠、高效传输都是基石。I2S(Inter-IC Sound)总线标准,正是为此而生的“行业通用语言”。它不像…

2026/7/26 1:18:01 阅读更多 →
Android 蓝牙开发详解与实战:基于 BluetoothChat 示例的完整解析

Android 蓝牙开发详解与实战:基于 BluetoothChat 示例的完整解析

摘要摘要:本文深入解析 Android SDK 中的经典蓝牙通信示例 BluetoothChat,从项目结构、核心代码到布局文件进行完整剖析。通过主界面 Activity、后台服务、设备列表、权限配置等模块的详细讲解,帮助开发者快速掌握 Android 蓝牙开发的核心流程…

2026/7/26 1:18:01 阅读更多 →
LoadRunner 自定义函数的使用方法

LoadRunner 自定义函数的使用方法

以前写程序时,不管是什么语言都支持函数。现在在学loadrunner ,不知道loadrunner是否支持自定义函数,于是做了个实验证明了一下,loadrunner也可以自定义函数。在下面的程序中,自定义函数max()放…

2026/7/26 1:18:01 阅读更多 →
3分钟上手:OBS AI背景移除插件完整使用指南

3分钟上手:OBS AI背景移除插件完整使用指南

3分钟上手:OBS AI背景移除插件完整使用指南 【免费下载链接】obs-backgroundremoval An OBS plugin for removing background in portrait images (video), making it easy to replace the background when recording or streaming. 项目地址: https://gitcode.co…

2026/7/26 1:17:01 阅读更多 →
卡美德生物科普:SAA1(血清淀粉样蛋白 A1)炎症相关靶点解析

卡美德生物科普:SAA1(血清淀粉样蛋白 A1)炎症相关靶点解析

SAA1 全称 Serum Amyloid A1,中文常称为血清淀粉样蛋白 A1,是机体炎症反应与免疫激活过程中具有代表性的分泌型蛋白分子,也是感染、炎症、组织损伤相关研究中常见的检测靶点。本文从分子基础、表达调控、检测意义与实验研究方向展开梳理&…

2026/7/26 1:17:01 阅读更多 →

日新闻

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 数据集6000张 完整源码已标注数据集训练好的模型环境配置教程程序运行说明文档,可以直接使用!系统支持图片、视频、摄像头等多种方式检测裂缝,功能强大实用。 1数据集6000张 8各类别

2026/7/26 0:00:31 阅读更多 →
深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

pubg数据集 精选原图1.42万数据 1.49万标签 无任何重复、算法增强或冗余图像! pubg绝地求生目标检测数据集 1分类:e_body,14905个标签,txt格式 共计14244张图,99%为640*640尺寸图像 适合yolo目标检测、AI训练关键词&am…

2026/7/26 0:00:31 阅读更多 →
Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex检测数据集数据集详情检测类别: allies enemy tag图片总量:7247张训练集:5139张验证集:1425张测试集:683张标注状态:全部已标注,即拿即用数据格式:支持YOLO格式及其他格式&#…

2026/7/26 0:00:31 阅读更多 →

周新闻

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 数据集6000张 完整源码已标注数据集训练好的模型环境配置教程程序运行说明文档,可以直接使用!系统支持图片、视频、摄像头等多种方式检测裂缝,功能强大实用。 1数据集6000张 8各类别

2026/7/26 0:00:31 阅读更多 →
深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

pubg数据集 精选原图1.42万数据 1.49万标签 无任何重复、算法增强或冗余图像! pubg绝地求生目标检测数据集 1分类:e_body,14905个标签,txt格式 共计14244张图,99%为640*640尺寸图像 适合yolo目标检测、AI训练关键词&am…

2026/7/26 0:00:31 阅读更多 →
Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex检测数据集数据集详情检测类别: allies enemy tag图片总量:7247张训练集:5139张验证集:1425张测试集:683张标注状态:全部已标注,即拿即用数据格式:支持YOLO格式及其他格式&#…

2026/7/26 0:00:31 阅读更多 →

月新闻