多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Redis List做消息队列的工程实践:从基础命令到可靠设计

Redis List做消息队列的工程实践:从基础命令到可靠设计 1. 项目起点为什么我在有MQ的时代还想用Redis List做队列先交代一下背景。我们有一个中小规模的业务系统用户量不大但业务链条却很长。比如用户提交一个导入任务后台要经过文件解析、数据清洗、格式校验、落库、生成报表通知五个步骤串行下来单次耗时可能超过十秒。一开始是同步调用接口超时被网关拦截用户前端转圈转到怀疑人生后来改成异步任务表加定时轮询任务多了调度器压力大数据库频繁被扫半夜还要被慢查询日志吵醒。当时也考虑过引入完整的消息队列中间件。但现实是团队就几个人运维精力有限机房资源也紧张为了一套独立MQ集群要维护Broker、NameServer、Consumer Rebalance机制投入产出比很低。后来我重新审视了项目里本来就有的基础设施发现Redis其实一直躺着没被充分利用。于是决定用Redis List作为消息队列的底层存储把任务流的异步化先跑起来。这就是这篇文章的由来。用通俗的话说Redis List就是一个双向链表你可以在头部左侧写入在尾部右侧弹出天然就像个水管——一头进水一头出水。最早接触Redis的开发者基本都写过这类代码但大多数人对它的认知停留在“一个能放列表的数据结构”很少有人把它真正当作消息队列的主力引擎去设计。本篇文章不讲虚的就讲这个List型消息队列能解决什么问题、怎么设计更靠谱、实际落地会遇到哪些坑。适合的人群很明确已经在业务里用Redis但不想为了异步任务专门引一套重量级MQ的团队以及那些做技术选型时还没搞清楚Redis List和专门MQ之间边界在哪的开发者。看完这篇文章你应该能判断自己的场景适不适合用List队列以及如果决定用了怎么做到尽量可靠。2. List队列的两条操作主线LPUSH与BRPOP的组合逻辑Redis List作为队列使用核心套路非常简单生产者往列表左侧塞消息消费者从右侧取消息一个进一个出天然先进先出。但就是这简单的LPUSH加RPOP背后隐藏着很多对效率和使用方式影响巨大的语义细节。2.1 阻塞式读取才是List队列的精髓很多初学者第一次写List队列消费端会用RPOP循环拉取。代码大概长这样import time import redis r redis.Redis(host127.0.0.1, port6379, db0) while True: msg r.rpop(task_queue) if msg is None: time.sleep(0.5) else: process(msg)这段代码的问题是极其明显的当队列为空时消费端在空转轮询每次RPOP都会发起一次完整的网络请求和Redis命令执行假如1秒轮询2次一个消费者一天下来要对Redis发起超过17万次无效查询多个消费者叠加就是一场小型网络风暴。正确做法是BRPOP。BRPOP的B是Blocking的缩写命令语义是如果队列里没有数据就阻塞等待直到有新的消息被LPUSH进来或者到达我们设定的超时时间。Redis作者在实现BRPOP时做了非常优雅的处理——当多个客户端同时阻塞等待同一个key时一旦有消息到达Redis会按照先来后到的顺序把消息分发给最早阻塞的客户端。while True: # 阻塞等待消息最多等30秒 msg r.brpop(task_queue, timeout30) if msg is None: continue # brpop返回二元组 (队列名, 消息内容) _, payload msg process(payload)这种改动的收益立竿见影消费端几乎不再发送无效请求消息一到立刻被取走延迟控制在毫秒级。而且BRPOP的超时时间并不是说30秒到了就一定返回空Redis的阻塞实现是基于事件驱动的如果30秒内有消息到达会立刻解除阻塞返回数据而不是“撑满30秒才响应”。2.2 多消费者场景下的消息分发特性再往里挖一层。如果你启动了两个消费端都用BRPOP监听同一个队列消息会怎么分配答案是每条消息只会被其中一个消费者拿到。Redis的BRPOP内部实现了类似互斥的语义——一个key被一个阻塞客户端取走数据后不会同时再发给其他阻塞客户端。这和Kafka里同一个消费组内部分区消息分配给不同消费者的逻辑是殊途同归的。但这里必须要说清楚一个边界Redis List队列天然支持的是竞争消费模式也就是消息只会被处理一次而不是Kafka那样的广播消费模式每条消息每个消费者都能处理一次。我见过不少团队想用List队列搞“一条消息同时触发邮件发送和短信通知”这搞错了模型广播场景应该用Redis的Pub/Sub或者直接在代码里让单个消费者分发到多个下游处理函数。如果确实需要“每条消息能被多个业务模块各自消费”有两种直觉但错误的方式经常出现每个模块都BRPOP同一个队列——每条消息被随机分给其中一个模块其他模块永远看不到这条消息。每个模块都LPUSH一份相同内容到各自队列——消费端代码变得很别扭生产者要感知所有下游模块。正确的解法应该是在消费者内部做路由分发取到消息后根据消息类型字段同时调用多个处理器或者使用Redis Stream的消费组特性。关于Redis Stream的对比后面专门展开。2.3 队列长度与内存的数学关系既然是用内存做存储就一定会关心容量问题。一个List能装多少消息完全取决于消息大小和Redis实例内存上限。我们来算一笔账。假设一条消息是200字节的JSON1000条消息大约是200KB100万条则是200MB。如果你的Redis实例内存上限是2GB这个队列大约能囤1000万条消息而不出问题。看起来数字很大但注意一个致命点Redis除了这个队列还有别的业务数据而且真正的内存占用还包括每条消息在List节点上的额外开销约几十字节不是单纯的数据大小。我在实际项目里给List队列设置了两个保护机制生产端写入前检查LLEN队列长度超过阈值比如5万条直接报警说明消费端已经跟不上了。Redis的maxmemory-policy设置成noeviction宁可让写入报错也不能让Redis把业务缓存数据给挤掉。# redis.conf 关键配置 maxmemory 2gb maxmemory-policy noeviction这个配置的含义是内存用满后新写入直接报错而不是悄悄淘汰旧数据。对队列场景来说消息丢失是不可接受的宁可让生产端报错触发重试也不能让队列里的消息被内存淘汰策略悄悄清掉。3. 消息队列选型对比List、Pub/Sub、Stream与专业MQ的边界写到这里一定有人会问Redis 5.0不是已经有Stream了吗它不是专门做消息队列的吗为什么还要用List还有Redis自己的Pub/Sub不也是消息模型吗这些疑问很合理我用一张表格把这几个方案的核心差异说清楚。维度ListPub/SubStream专业MQ(RabbitMQ/Kafka)消息持久化有RDB/AOF无实时推送有可持久化有不同存储策略消息堆积内存存储受限于maxmemory不堆积实时丢弃可写盘堆积能力强磁盘存储堆积能力强消费确认无专门ACK需自行设计无有消费组ACK机制完善的ACK与重投重复消费不支持弹出即消失不支持支持从指定ID重读按offset/replay支持复杂度极简简单中等高需部署维护典型场景轻量任务异步化实时广播通知需要ACK和追溯的通用队列大型分布式核心链路从表格能看出来List队列的本质是一个“弹出去就没了”的模型一旦RPOP成功这条消息就从Redis里消失了。这带来了两个直接问题一是如果消费者处理消息时崩溃这条消息就找不回来了二是无法回溯历史消息做排查。这是List队列的短板但不是不能通过设计来规避后面章节会专门讲。那Stream出来后List还有没有存在价值我的观点是List虽然朴素但它足够简单没有消费组、ACK、Pending Entries这些概念心智负担低。而且BRPOP阻塞语义比XREAD的阻塞参数更直观。如果你的业务只需要“把任务塞进队列消费者取走处理”这一个动作不需要确认机制List的简洁性反而是一种优势——代码短不容易出错排障参数少。Pub/Sub则完全是两个物种。它不存储消息生产端PUBLISH后消息立即推给所有订阅者如果此刻没有订阅者在线消息直接丢失。用Pub/Sub做任务队列是新手最容易踩的坑——发布任务时消费者恰好重启任务就凭空消失。Pub/Sub适合做在线广播比如配置变更通知、刷新本地缓存的指令不适合做任务队列。Stream更像是List的进化版它保留了List那种追加消息的模型同时加入了消费组、消费者ID、ACK机制本质上就是内置了消费者管理和消息确认的专业队列。如果你的项目已经用上Redis 5.0而且对消息可靠性有更高要求我建议直接评估Stream但如果Redis版本是3.x、4.x这种老版本List绝对是最务实的选择。还有一点值得提醒从运维角度另一个容易被忽略的选择是我是否应该直接用云厂商提供的消息队列服务。如果公司买的是云数据库Redis那List队列仍然共享同一个实例的连接数和内存资源队列规模一大就会影响其他业务。这个容量边界要提前跟DBA沟通清楚预留独立的Redis实例专门做队列避免互相干扰。4. 实战演进从简单队列到带确认、重试、延迟的可靠队列进入真正落地环节。这一节我按照从简到繁的顺序展示一个生产可用的List消息队列怎么一步步搭起来。很多网上的示例只做到了“LPUSHBRPOP”这一步但真实业务根本没有那么简单——任务处理失败怎么办、消费者崩溃消息丢失怎么办、定时任务和延时任务怎么做这些都是必须正面回答的问题。4.1 基础版单一队列生产与消费首先是标准的生产者消费者模型。以Python为例演示但思路各语言通用。生产者import json import redis r redis.Redis(host127.0.0.1, port6379, db0) def publish(queue_name, task_data): message json.dumps(task_data, ensure_asciiFalse) r.lpush(queue_name, message) return len(message) # 可以返回长度做日志统计 # 实际调用示例 publish(import_task_queue, { task_id: T20250115001, type: file_import, file_url: https://oss.example.com/upload/file.csv, user_id: 1024, created_at: 2025-01-15 10:30:00 })消费者import json import redis r redis.Redis(host127.0.0.1, port6379, db0) def process_message(message): task json.loads(message) # 模拟业务处理 print(f处理任务: {task[task_id]}, 类型: {task[type]}) # 这里放真正的业务逻辑 while True: result r.brpop([import_task_queue, default_queue], timeout30) if result is None: continue queue_name, message result try: process_message(message) except Exception as exc: print(f任务处理异常: {exc})这里有个小细节BRPOP的key参数可以传多个队列名Redis会从左到右依次检查哪个队列有数据优先返回最早有数据的那个。这个特性可以做优先级队列把紧急任务的队列名放在前面普通队列放后面。4.2 进阶版如何用备份队列和ACK实现“伪确认”前面说了List弹出即消失为了弥补这个缺陷我采用的方案是“先处理、后确认”加“失败丢回备份队列”。所谓“先处理、后确认”是指消费者BRPOP拿到消息后先不急着从语义上“消费掉”这条消息——因为BRPOP已经把它从主队列里取走了所以只能在业务层面做文章。我的做法是消费者BRPOP拿到消息。立刻将消息内容LPUSH到对应的“处理中队列”(queue_name:processing)相当于是先把消息放到一个临时暂存区。执行业务逻辑。业务成功从处理中队列删除这条消息LREM按值删除。业务失败把消息从处理中队列弹出来塞回主队列或者塞入延迟重试队列。这段逻辑用伪代码表示import json import redis import time r redis.Redis(host127.0.0.1, port6379, db0) def process_with_ack(queue_name, message, max_retry3): # 进入处理中状态 processing_key f{queue_name}:processing r.lpush(processing_key, message) task json.loads(message) try: # 业务处理假设这里可能抛异常 do_business(task) # 成功后从处理中队列移除 r.lrem(processing_key, 1, message) return True except Exception as exc: # 检查重试次数是否超限 retry_key f{queue_name}:retry:{task[task_id]} retry_count r.incr(retry_key) r.expire(retry_key, 86400) # 24小时后自动过期 if retry_count max_retry: # 放回队列稍后重试 r.rpush(queue_name, message) else: # 超过重试次数进入死信队列 r.lpush(f{queue_name}:dead, message) # 清理处理中记录 r.lrem(processing_key, 1, message) return False这个方案并不能百分百保证不丢消息——如果消费者在处理过程中直接宕机消息会一直残留在processing队列里永远不会被重新投递。真正严谨的做法是给处理中消息设置过期扫描任务定期从processing队列里捞出滞留超过时间阈值的消息重新投递。这个“孤儿消息”扫描不是简单的遍历因为LREM只能按值删按值找效率会随队列长度退化更可靠的做法是用Stream。但对于中小业务我上面这段代码配合Redis的过期时间已经能把消息丢失率降到很低的水平。4.3 进阶版延迟队列和优先级队列的实现思路业务里经常有“30分钟后未支付关闭订单”这类延迟任务需求。用List实现延迟队列本质上是用时间排序的ZSet加一个轮询线程。思路是这样任务到期时间作为Score存入ZSet。单独起一个定时轮询线程周期性地用ZRANGEBYSCORE查出所有Score小于等于当前时间的任务。把这些任务从ZSet移除ZREM再LPUSH到真正的业务队列。import time import redis r redis.Redis(host127.0.0.1, port6379, db0) DELAY_QUEUE_KEY delay:order_close def add_delay_task(task_id, delay_seconds): 把任务加入延迟队列delay_seconds后可见 score time.time() delay_seconds r.zadd(DELAY_QUEUE_KEY, {task_id: score}) def poll_delay_queue(): 轮询把到期的延迟任务转移到业务队列 now time.time() # 取出所有到期的任务ID expired r.zrangebyscore(DELAY_QUEUE_KEY, 0, now) if expired: # 从延迟队列移除同时加入业务队列 pipeline r.pipeline() for task_id in expired: pipeline.zrem(DELAY_QUEUE_KEY, task_id) pipeline.lpush(order_close_queue, task_id) pipeline.execute() # 启动一个后台线程每5秒执行一次 import threading def start_polling(): while True: poll_delay_queue() time.sleep(5) t threading.Thread(targetstart_polling, daemonTrue) t.start()优先级队列的核心在前面的BRPOP多队列参数里已经提到过一次。Redis官方没有提供权衡多个队列优先级的原语但我们可以通过BRPOP key1 key2的入参顺序让高优先级队列先被消费。举个例子BRPOP先监听urgent_task_queue再监听normal_task_queue那么只要有紧急消息消费者会优先取走它普通消息如果没紧急消息才会被取到。不过要说清楚这只是一个近似优先级不是严格的优先级调度。如果高优先级队列持续有消息低优先级队列会一直饥饿。所以我的建议是区分出“紧急告警”和“普通任务”两级就够了不要试图用它实现多级复杂的调度那是专业MQ的活。4.4 进阶版批量消费提升吞吐量List队列的消息是一条条弹出的如果业务是那种单条处理代价大但批量处理效率高的类型比如批量写数据库、批量调接口就需要让消费者攒一批消息再处理。Redis没有原生的批量BRPOP但可以用RPOPLPUSH组合和循环多次RPOP实现。简单做法是用非阻塞的LRANGE批量取出加LTRIM裁剪或者多次RPOPdef batch_pop(queue_name, batch_size10): 批量取出消息最多取batch_size条 # 使用pipeline减少网络往返 pipeline r.pipeline() for _ in range(batch_size): pipeline.rpop(queue_name) results pipeline.execute() return [msg for msg in results if msg is not None]这样有什么问题如果队列里消息数量少于batch_size该轮只能取到部分消息批量处理的批次不够大效率提升有限。更高效的做法是让生产者生产一条总控消息然后消费者根据总控消息去批量取一批子消息处理。这种模式适合“批量导入、批量导出”类任务我在实际项目里就是这样设计的——总控消息里包含一个批次ID消费者拿到后根据批次ID把所有关联子任务从List中批量取出一次性处理效率比逐条消费高了一个数量级。5. 重复消费与消息丢失排查链路与设计防线List队列用久了一定会遇到两个高频问题消息重复消费了或者消息莫名丢失了。这一节我把这两类问题各自的排查思路完整写出来这些都是实打实踩出来的经验。5.1 消息重复消费的根因与业务幂等设计List队列本身不会有重复消费的问题——BRPOP取出后消息就没了Redis不会像Kafka那样由于offset提交失败造成消息重复投递。但为什么实际运行时还是会出现重复处理我遇到过的根因有三类第一类是网络超时导致的重试。生产者LPUSH后网络恰好断开客户端无法确定Redis是否收到了消息于是重试LPUSH同一业务消息被塞进了队列两次。这类问题可以在生产端维护一个消息唯一ID连同业务数据一起写入队列。第二类是消费者端语义导致的重复。BRPOP拿到消息后业务处理超时监控系统误判消费者死了于是重启消费者但此时消息已经被取下正在处理重启后消息就丢了而不是重复了。这个严格说不是重复消费而是消息丢失的变种。第三类是人为重放。排查问题时手动把一条旧消息重新塞回队列忘记摘除业务状态导致下游重复执行。把这三类根因串起来后结论清晰即使List本身不会主动重复但分布式环境的每个环节都可能导致重复语义。最稳妥的防线不是依赖队列而是业务层幂等。“幂等”的意思是同一条消息无论处理多少次最终结果都和只处理一次一致。幂等实现最常用的方案是唯一ID加去重表。以订单发货为例消息里带上order_id和event_id每次处理前先去Redis查这个event_id是否处理过def is_duplicate_message(event_id): # 如果之前没处理过设置一个10分钟的去重标记 return not r.set(fmsg:dup:{event_id}, 1, nxTrue, ex600)这个SET命令带上NX参数表示“只有key不存在时才设置成功”如果返回False说明这条消息已经被处理过直接跳过。这个方案实现简单IO开销小是List队列最推荐的幂等防线。5.2 消息丢失的链路排查从生产端到消费端逐个击破对比重复消费消息丢失更容易被忽视因为丢失时往往没有任何报错。排查时我建议按以下链路逐段分析第一段生产端是否真的成功写入。很多开发者写代码时忽略了LPUSH的返回值网络抖动、连接池耗尽、内存淘汰都可能导致LPUSH失败但代码继续往下走。务必检查LPUSH的返回值如果返回0或抛异常要在业务层面记录日志并触发告警。第二段Redis内存淘汰策略是否吞掉了消息。如果Redis实例的maxmemory-policy设置成了allkeys-lru或volatile-lru内存紧张时List里的消息会被当作冷数据淘汰消息悄无声息地消失。这也是我在2.3节强调必须设置noeviction的根本原因。第三段消费者是否真的处理成功。如果BRPOP取到消息后由于业务逻辑异常导致任务中断消息就已经从队列中消失了没有重试机制的话就彻底丢了。需要引入4.2节讲的处理中队列机制来兜底。把三段串起来我给出一个完整可靠的消息队列实现已经在一套日吞吐数十万任务的系统上稳定运行过环节措施生产端幂等ID、检查LPUSH返回值、失败告警存储端maxmemory-policynoeviction、独立实例、LLEN监控消费端处理中队列、失败重试、死信队列、幂等消费兜底定时扫描processing队列中的滞留消息5.3 死信队列处理那些“永远处理不了”的消息死信队列是消息系统里一个容易被忽视但极端重要的设计。所谓死信就是重试了N次依然处理失败的消息。如果这些消息一直留在主队列里反复重试不仅浪费消费端资源还会阻塞后续正常消息甚至导致整个队列消费停滞。我在代码里加入了一个死信队列机制每条消息记录重试次数重试超过阈值就把它从主队列移动到xxx:dead队列同时发出告警。定时任务每天扫描死信队列人工处理或自动丢弃。还需要给每条消息记录它的失败原因。List无法存储元数据所以我在消息体内塞了一个retry_count字段每次失败重试前更新它。当然这只是List的变通方案Stream可以用MessageId和Pending机制做得更优雅。6. 当队列变大了内存监控、性能调优与容量规划队列在开发环境跑得再欢到了生产环境一压流量就原形毕露。最后这部分聊聊队列规模化之后的性能问题和容量规划。6.1 队列长度与消费者数量的匹配策略很多人问“我是不是消费者开得越多消费速度就越快”这个问题的答案是不一定。List队列的消费速度受限于同一时刻BRPOP消息获取的串行语义和下游业务处理能力。Redis本身单线程处理命令它的性能瓶颈通常在CPU单核而不是并发连接数。对于List队列Redis吞吐量能到每秒数万甚至十万级别的LPUSH/BRPOP操作但你的消费者如果每个都要做复杂的业务计算瓶颈就在业务侧而不是队列侧。一个基本的估算公式单个消费者处理能力 1000 / 单条消息平均处理耗时(毫秒) 需要消费者数量 目标QPS / 单个消费者处理能力举个例子如果单条消息处理耗时200ms单个消费者的处理能力就是5 QPS。如果目标QPS是50就需要至少10个消费者实例。注意这是理想情况实际还会被下游数据库连接数、网络IO等因素限制所以建议留出30%到50%的余量。6.2 内存增长的告警与消费积压监控List队列不同于专业MQ消息堆积会直接占用Redis内存。所以监控必不可少。我在生产环境里有一套组合拳用LLEN监控每个队列长度接入Prometheus定时采集队列长度超过预设阈值就触发告警。监控Redis实例的used_memory设定基于maxmemory的告警线比如达到80%就预警。消费者端记录消费延迟——一条消息从LPUSH到BRPOP取出的时间差这个指标能直接反映消费堆积的严重程度。消费延迟可以用消息体里带一个created_at字段消费者处理时计算当前时间减去它。如果延迟从秒级涨到分钟级说明消费者已经忙不过来了。6.3 AOF与主从复制对队列可靠性的影响这里说一个很多人忽略的点Redis的持久化策略对List队列的影响。RDB快照是定时打点的如果Redis在两次快照之间宕机就会丢失最近几分钟的数据AOF是追加写日志可以配置appendfsync everysec或者always丢失窗口更小。对于重要的消息队列场景建议至少开启AOF并且设置appendfsync everysec这样最多只会丢1秒内的数据。如果业务对丢消息零容忍那还是要放弃List队列考虑Stream或者专业MQ。另外Redis主从复制也是保障队列可靠性的重要手段。主节点宕机后从节点可以顶上。但必须注意Redis的主从复制是异步的极端情况下主节点写入成功后还没把数据同步给从节点就宕机消息还是会丢。所以主从复制只能提高可用性不能解决数据丢失问题。7. 最终的技术判断哪些场景适合List队列哪些该用专门MQ分享一点个人的判断标准。我在多个项目里反复比较过总结出几个能直接套用的结论。适合用Redis List队列的场景业务日均消息量在百万级以下单个消息不超过几十KB。对消息可靠性的要求是“最多能忍受秒级或分钟级的少量丢失”但可以通过业务幂等兜底。团队运维能力有限不想再维护一套中间件集群。项目已经引入Redis通过复用基础设施降低成本。场景简单就是“一个任务排队一个消费者处理”不需要复杂路由、重投策略。不适合用Redis List队列的场景金融交易、订单支付这类对消息零丢失要求的核心链路。需要根据消息内容做复杂路由分发的场景。海量消息堆积超出Redis内存承载上限的场景。需要消息回溯、审计追查的历史事件流场景。我自己在项目里的实践是List队列用在导入导出、报表生成、通知推送这类弱一致性任务上非常顺手成本几乎为零出问题时通过业务幂等就能补救。一旦涉及资金、状态流转等核心链路我会直接交给专业MQ。如果你正处于“要不要用Redis List做消息队列”的犹豫阶段我的建议很直接先想清楚业务对消息丢失的容忍度再想清楚团队运维能力然后做技术选型。别被“高性能”“高并发”这种词带偏消息队列的核心从来不是性能而是可靠性模型。List能提供的可靠性上限就是“内存里躺着就不会丢取走了就看命”的弱保证用合理的业务设计补足之后它完全能撑起一个中小规模系统的异步化要求。这套方案我用了很久踩过坑也填过坑写出来给后来人做个参考。如果你已经在用Redis List做队列欢迎在评论区聊聊你遇到过的坑互相补全一下边界条件。
返回列表