
延迟队列(Delay Queue)是分布式系统中极为常见的基础设施。无论是订单超时取消、消息定时推送、重试退避调度,还是限流令牌发放,核心诉求都一样:把一个任务延迟到指定时间点后再执行。市面上有 RabbitMQ 死信队列、Kafka 时间轮、RocketMQ 延迟消息等多种方案,但如果你的技术栈中已经有 Redis,那么用 Sorted Set(ZSet)实现延迟队列往往是最轻量、最可控的选择。
本文将深入剖析 Redis ZSet 延迟队列的实现原理,给出完整的 Python 生产级代码,并系统性地讨论容灾、幂等、并发控制等生产环境必须面对的问题。
一、延迟队列的核心应用场景
在动手实现之前,先明确延迟队列到底解决什么问题。以下是几个典型的业务场景:
- 订单超时关闭:用户下单后 30 分钟未支付,自动关闭订单并释放库存。这是电商系统最高频的延迟任务。
- 消息定时推送:用户预约了早上 9 点的会议提醒,系统需要在指定时间点发送推送通知。
- 异步重试退避:调用第三方 API 失败后,按照指数退避策略(1s、2s、4s、8s…)延迟重试,避免雪崩。
- 限流令牌补充:令牌桶限流器中,消费令牌后按固定速率延迟补充。
- 心跳超时检测:分布式系统中检测节点是否存活,若 N 秒内未收到心跳则标记为下线。
这些场景的共同特征是:任务需要在未来的某个精确时间点被触发,且必须保证可靠性——即使服务重启也不能丢失。
二、常见延迟队列方案对比
在选择 Redis ZSet 方案之前,有必要了解主流方案的优劣,以便做出合理的技术选型:
| 方案 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| Redis ZSet | 实现简单、精度高、支持任意延迟时间、与现有 Redis 基础设施复用 | 单实例吞吐有限、无原生 ACK 机制、需自行处理容灾 | 中小规模延迟任务、已有 Redis 基础设施 |
| RabbitMQ TTL+DLX | 成熟稳定、原生 ACK、支持持久化 | 配置复杂、TTL 存在头部阻塞问题、延迟精度受消息顺序影响 | 已使用 RabbitMQ 的系统、需要强可靠性 |
| RocketMQ 延迟消息 | 原生支持延迟级别、高吞吐、持久化 | 开源版仅支持固定延迟级别(1s/5s/10s/…)、不可任意延迟 | 大规模消息系统、固定延迟级别场景 |
| Kafka + 时间轮 | 超高吞吐、持久化、可回溯 | 实现复杂、延迟精度依赖轮询间隔、无原生延迟语义 | 超大规模流处理、日志场景 |
| Java 时间轮(Netty HashedWheelTimer) | 纯内存、极高性能、无外部依赖 | 单机方案、服务重启丢数据、无持久化 | 单机高频延迟任务、可容忍丢失 |
从对比可以看出,Redis ZSet 方案在「已有 Redis 基础设施 + 中小规模延迟任务」的场景下,是性价比最高的选择。下面我们进入实现部分。
三、Redis ZSet 实现延迟队列的核心原理
3.1 基本原理
Redis Sorted Set 的核心特性是:每个成员关联一个 score,集合按 score 排序。如果把 延迟任务的执行时间戳 作为 score,把 任务 ID 或任务数据 作为 member,就天然得到了一个按时间排序的延迟队列。
核心操作流程如下:
- 生产者:将任务以
1ZADD delay_queue <timestamp> <task_id>
写入 ZSet,score 为任务应该执行的时间戳(毫秒)。
- 消费者:定时轮询
1ZRANGEBYSCORE delay_queue 0 <current_timestamp> LIMIT 0 N
,获取所有已到期的任务。
- 消费成功:从 ZSet 中移除已消费的任务
1ZREM delay_queue <task_id>
。
- 消费失败:将任务重新 ZADD 回队列,score 设为下一次重试时间。
3.2 为什么选择 ZSet 而不是 List
有人可能想到用 Redis List + LPUSH/RPOP 实现延迟队列,但 List 是 FIFO 结构,无法按时间排序。如果先入队的任务延迟时间更长,后入队的任务延迟时间更短,List 就无法正确调度。ZSet 的 score 排序特性恰好解决了这个问题。
3.3 关键命令清单
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17 # 添加延迟任务(score = 执行时间戳,member = 任务ID)
ZADD delay_queue 1695000000000 "task:order:12345"
# 查询已到期任务(当前时间戳之前的所有任务,取前100条)
ZRANGEBYSCORE delay_queue 0 1695000000000 LIMIT 0 100
# 移除已消费任务
ZREM delay_queue "task:order:12345"
# 查看队列大小
ZCARD delay_queue
# 查看最近一条任务的执行时间
ZRANGE delay_queue 0 0 WITHSCORES
# 清理过期任务(按score范围删除)
ZREMRANGEBYSCORE delay_queue 0 1694999999999
四、生产级 Python 代码实现
下面给出完整的 Python 实现,包含生产者和消费者两个模块。使用
1 | redis-py |
库,兼容 Redis Cluster。
4.1 延迟队列核心类
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161 import json
import time
import uuid
import logging
import redis
from typing import Optional, Dict, Any
from dataclasses import dataclass, asdict
logger = logging.getLogger(__name__)
@dataclass
class DelayTask:
"""延迟任务数据结构"""
task_id: str
task_type: str # 任务类型,如 order_timeout, retry_api
payload: Dict[str, Any] # 任务负载
max_retries: int = 3 # 最大重试次数
retry_count: int = 0 # 当前重试次数
def to_json(self) -> str:
return json.dumps(asdict(self), ensure_ascii=False)
@classmethod
def from_json(cls, data: str) -> 'DelayTask':
d = json.loads(data)
return cls(**d)
class RedisDelayQueue:
"""Redis ZSet 延迟队列"""
def __init__(self, redis_client: redis.Redis, queue_key: str = "delay_queue"):
self.r = redis_client
self.queue_key = queue_key
# 用 Hash 存储任务详情,ZSet 只存 task_id 做时间索引
self.task_data_key = f"{queue_key}:data"
def add_task(self, task: DelayTask, delay_seconds: float) -> bool:
"""添加延迟任务
Args:
task: 任务对象
delay_seconds: 延迟秒数(从此刻开始计算)
Returns:
True if added successfully
"""
execute_at = int(time.time() * 1000 + delay_seconds * 1000)
pipe = self.r.pipeline()
# 存任务详情到 Hash
pipe.hset(self.task_data_key, task.task_id, task.to_json())
# 存时间索引到 ZSet
pipe.zadd(self.queue_key, {task.task_id: execute_at})
pipe.execute()
logger.info(f"Task added: {task.task_id}, execute_at={execute_at}")
return True
def add_task_at(self, task: DelayTask, execute_at_timestamp_ms: int) -> bool:
"""添加任务到指定时间点执行(绝对时间戳,毫秒)"""
pipe = self.r.pipeline()
pipe.hset(self.task_data_key, task.task_id, task.to_json())
pipe.zadd(self.queue_key, {task.task_id: execute_at_timestamp_ms})
pipe.execute()
return True
def fetch_ready_tasks(self, batch_size: int = 100) -> list:
"""获取已到期的任务(使用 Lua 脚本保证原子性)
使用 ZPOPMIN 原子弹出,避免多消费者竞争问题
"""
now_ms = int(time.time() * 1000)
# Lua 脚本:原子地取出到期任务
lua_script = """
local now = ARGV[1]
local limit = ARGV[2]
local queue_key = KEYS[1]
local data_key = KEYS[2]
-- 取出到期任务ID
local task_ids = redis.call('ZRANGEBYSCORE', queue_key, 0, now, 'LIMIT', 0, limit)
if #task_ids == 0 then
return {}
end
-- 取出任务详情
local results = {}
for i, task_id in ipairs(task_ids) do
local data = redis.call('HGET', data_key, task_id)
if data then
table.insert(results, task_id)
table.insert(results, data)
end
end
-- 从ZSet移除已取出的任务
for i, task_id in ipairs(task_ids) do
redis.call('ZREM', queue_key, task_id)
end
return results
"""
raw = self.r.eval(lua_script, 2, self.queue_key, self.task_data_key, now_ms, batch_size)
tasks = []
for i in range(0, len(raw), 2):
task_id = raw[i].decode() if isinstance(raw[i], bytes) else raw[i]
task_data = raw[i+1].decode() if isinstance(raw[i+1], bytes) else raw[i+1]
tasks.append((task_id, DelayTask.from_json(task_data)))
if tasks:
logger.info(f"Fetched {len(tasks)} ready tasks")
return tasks
def ack_task(self, task_id: str) -> bool:
"""确认任务完成,清理任务数据"""
pipe = self.r.pipeline()
pipe.zrem(self.queue_key, task_id)
pipe.hdel(self.task_data_key, task_id)
pipe.execute()
return True
def nack_task(self, task_id: str, task: DelayTask, retry_delay: float = 5.0) -> bool:
"""任务处理失败,重新入队
Args:
task_id: 任务ID
task: 任务对象(retry_count会自增)
retry_delay: 重试延迟秒数
"""
task.retry_count += 1
if task.retry_count > task.max_retries:
# 超过最大重试次数,放入死信队列
self._move_to_dlq(task_id, task)
logger.warning(f"Task {task_id} moved to DLQ after {task.retry_count} retries")
return False
# 指数退避重试
backoff_delay = retry_delay * (2 ** (task.retry_count - 1))
execute_at = int(time.time() * 1000 + backoff_delay * 1000)
pipe = self.r.pipeline()
pipe.hset(self.task_data_key, task_id, task.to_json())
pipe.zadd(self.queue_key, {task_id: execute_at})
pipe.execute()
logger.info(f"Task {task_id} requeued, retry={task.retry_count}, delay={backoff_delay}s")
return True
def _move_to_dlq(self, task_id: str, task: DelayTask):
"""移入死信队列"""
dlq_key = f"{self.queue_key}:dlq"
self.r.hset(dlq_key, task_id, task.to_json())
self.r.zrem(self.queue_key, task_id)
self.r.hdel(self.task_data_key, task_id)
def queue_size(self) -> int:
return self.r.zcard(self.queue_key)
def next_task_time(self) -> Optional[int]:
"""返回下一条任务的执行时间戳(毫秒),用于计算轮询间隔"""
result = self.r.zrange(self.queue_key, 0, 0, withscores=True)
if result:
return int(result[0][1])
return None
4.2 消费者实现
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80 import threading
import time
import signal
import logging
from concurrent.futures import ThreadPoolExecutor
logger = logging.getLogger(__name__)
class DelayQueueConsumer:
"""延迟队列消费者"""
def __init__(self, queue: RedisDelayQueue, handler_map: dict,
poll_interval: float = 1.0, worker_count: int = 4):
self.queue = queue
self.handler_map = handler_map # {task_type: handler_func}
self.poll_interval = poll_interval
self.worker_count = worker_count
self._running = False
self._executor = ThreadPoolExecutor(max_workers=worker_count)
self._lock = threading.Lock()
def start(self):
self._running = True
signal.signal(signal.SIGTERM, self._handle_signal)
signal.signal(signal.SIGINT, self._handle_signal)
logger.info("Consumer started, poll_interval=%.2fs, workers=%d",
self.poll_interval, self.worker_count)
self._poll_loop()
def _poll_loop(self):
while self._running:
try:
tasks = self.queue.fetch_ready_tasks(batch_size=50)
if tasks:
for task_id, task in tasks:
self._executor.submit(self._process_task, task_id, task)
else:
# 动态调整轮询间隔:如果有下一条任务,按需等待
next_time = self.queue.next_task_time()
if next_time:
now_ms = int(time.time() * 1000)
sleep_time = max(0.1, min(
(next_time - now_ms) / 1000.0,
self.poll_interval
))
time.sleep(sleep_time)
else:
time.sleep(self.poll_interval)
except Exception as e:
logger.error(f"Poll loop error: {e}", exc_info=True)
time.sleep(self.poll_interval)
def _process_task(self, task_id: str, task: DelayTask):
"""处理单个任务"""
handler = self.handler_map.get(task.task_type)
if not handler:
logger.error(f"No handler for task_type={task.task_type}, task_id={task_id}")
self.queue.nack_task(task_id, task, retry_delay=60)
return
try:
result = handler(task.payload)
if result:
self.queue.ack_task(task_id)
logger.info(f"Task {task_id} processed successfully")
else:
self.queue.nack_task(task_id, task)
except Exception as e:
logger.error(f"Task {task_id} failed: {e}", exc_info=True)
self.queue.nack_task(task_id, task)
def _handle_signal(self, signum, frame):
logger.info(f"Received signal {signum}, shutting down...")
self._running = False
self._executor.shutdown(wait=True, timeout=30)
logger.info("Consumer stopped")
def stop(self):
self._running = False
self._executor.shutdown(wait=True)
4.3 完整使用示例
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62 import redis
import logging
logging.basicConfig(level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s')
# 初始化
r = redis.Redis(host='127.0.0.1', port=6379, db=0, decode_responses=False)
queue = RedisDelayQueue(r, queue_key="delay_queue:orders")
# ====== 定义任务处理器 ======
def handle_order_timeout(payload: dict) -> bool:
"""订单超时处理"""
order_id = payload.get('order_id')
logger.info(f"Processing order timeout: {order_id}")
# 实际业务:关闭订单、释放库存
# close_order(order_id)
# release_stock(order_id)
return True
def handle_api_retry(payload: dict) -> bool:
"""API重试任务"""
url = payload.get('url')
data = payload.get('data')
logger.info(f"Retrying API call to {url}")
# response = requests.post(url, json=data, timeout=5)
# if response.status_code == 200:
# return True
# return False
return True
handler_map = {
'order_timeout': handle_order_timeout,
'api_retry': handle_api_retry,
}
# ====== 生产者:添加延迟任务 ======
# 场景1:订单30分钟后超时关闭
task1 = DelayTask(
task_id=f"task:order:{uuid.uuid4().hex}",
task_type='order_timeout',
payload={'order_id': 'ORD-2026-001234', 'amount': 99.50}
)
queue.add_task(task1, delay_seconds=30 * 60) # 30分钟后执行
# 场景2:3秒后重试API调用
task2 = DelayTask(
task_id=f"task:retry:{uuid.uuid4().hex}",
task_type='api_retry',
payload={'url': 'https://api.example.com/webhook', 'data': {'event': 'payment'}}
)
queue.add_task(task2, delay_seconds=3)
# ====== 消费者:启动消费 ======
consumer = DelayQueueConsumer(
queue=queue,
handler_map=handler_map,
poll_interval=1.0,
worker_count=4
)
consumer.start()

五、生产环境必须解决的五个问题
5.1 并发消费的原子性:Lua 脚本
多个消费者实例同时轮询时,如果不做原子性控制,可能出现同一个任务被多个消费者取走的问题。上面的代码通过 Lua 脚本解决了这个问题——
1 | ZRANGEBYSCORE |
+
1 | ZREM |
在 Lua 脚本中作为一个原子操作执行,Redis 保证脚本执行期间不会被其他命令打断。
如果不用 Lua 脚本,也可以用
1 | ZPOPMIN |
命令,但它会弹出 score 最小的成员,不一定是已到期的。因此
1 | ZRANGEBYSCORE + ZREM |
的 Lua 组合是更精确的方案。
5.2 服务重启后的任务恢复
这是最容易忽视的问题。如果消费者在
1 | fetch_ready_tasks |
之后、
1 | ack_task |
之前崩溃,这个任务就「丢失」了——它已经从 ZSet 中移除,但还没被处理。
解决方案是引入 处理中状态(In-Flight):
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81 class RedisDelayQueue:
# ... 前面的代码不变 ...
def fetch_ready_tasks_safe(self, batch_size: int = 100) -> list:
"""安全版:取出任务时移入 processing 集合,设置超时"""
now_ms = int(time.time() * 1000)
processing_key = f"{self.queue_key}:processing"
processing_timeout = 60 # 60秒未ACK视为超时
lua_script = """
local now = ARGV[1]
local limit = ARGV[2]
local timeout = ARGV[3]
local queue_key = KEYS[1]
local data_key = KEYS[2]
local processing_key = KEYS[3]
local task_ids = redis.call('ZRANGEBYSCORE', queue_key, 0, now, 'LIMIT', 0, limit)
if #task_ids == 0 then
return {}
end
local results = {}
for i, task_id in ipairs(task_ids) do
local data = redis.call('HGET', data_key, task_id)
if data then
-- 从队列移除,加入processing集合,设置超时
redis.call('ZREM', queue_key, task_id)
redis.call('HSET', processing_key, task_id, data)
redis.call('ZADD', processing_key .. ':timeout', task_id, now + timeout * 1000)
table.insert(results, task_id)
table.insert(results, data)
end
end
return results
"""
raw = self.r.eval(lua_script, 3,
self.queue_key, self.task_data_key, processing_key,
now_ms, batch_size, processing_timeout)
tasks = []
for i in range(0, len(raw), 2):
task_id = raw[i].decode() if isinstance(raw[i], bytes) else raw[i]
task_data = raw[i+1].decode() if isinstance(raw[i+1], bytes) else raw[i+1]
tasks.append((task_id, DelayTask.from_json(task_data)))
return tasks
def ack_task_safe(self, task_id: str) -> bool:
"""安全版ACK:从processing集合清理"""
processing_key = f"{self.queue_key}:processing"
timeout_key = f"{processing_key}:timeout"
pipe = self.r.pipeline()
pipe.hdel(processing_key, task_id)
pipe.zrem(timeout_key, task_id)
pipe.hdel(self.task_data_key, task_id)
pipe.execute()
return True
def recover_stale_tasks(self):
"""恢复超时未ACK的任务(定时调用)"""
now_ms = int(time.time() * 1000)
processing_key = f"{self.queue_key}:processing"
timeout_key = f"{processing_key}:timeout"
# 找出超时任务
stale_ids = self.r.zrangebyscore(timeout_key, 0, now_ms)
recovered = 0
for task_id_bytes in stale_ids:
task_id = task_id_bytes.decode() if isinstance(task_id_bytes, bytes) else task_id_bytes
task_data = self.r.hget(processing_key, task_id)
if task_data:
task_data_str = task_data.decode() if isinstance(task_data, bytes) else task_data
task = DelayTask.from_json(task_data_str)
# 重新入队,延迟5秒
self.add_task(task, delay_seconds=5)
# 清理processing
self.r.hdel(processing_key, task_id)
self.r.zrem(timeout_key, task_id)
recovered += 1
if recovered > 0:
logger.info(f"Recovered {recovered} stale tasks")
return recovered
然后在消费者的轮询循环中定期调用
1 | recover_stale_tasks |
:
1
2
3
4
5
6
7
8
9
10 # 每30秒执行一次恢复
last_recover_time = 0
RECOVER_INTERVAL = 30
while self._running:
now = time.time()
if now - last_recover_time >= RECOVER_INTERVAL:
self.queue.recover_stale_tasks()
last_recover_time = now
# ... 正常消费逻辑 ...
5.3 幂等性保证
即使有了 processing 状态和恢复机制,任务仍可能被执行多次(比如恢复后再次被消费)。因此 业务处理器本身必须实现幂等性。推荐的做法是:
- 使用唯一任务 ID 作为幂等键:在执行业务操作前,先检查该 task_id 是否已处理过(可以用 Redis SETNX 标记)。
- 数据库唯一约束:在数据库层面对关键操作加唯一索引。
- 状态机检查:执行前检查业务对象的状态,只有处于预期状态才执行。
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18 def handle_order_timeout(payload: dict, task_id: str) -> bool:
"""幂等的订单超时处理"""
# 幂等检查1:Redis标记
idempotent_key = f"idempotent:{task_id}"
if not r.set(idempotent_key, '1', nx=True, ex=86400):
logger.info(f"Task {task_id} already processed, skip")
return True
# 幂等检查2:状态机
order = get_order(payload['order_id'])
if order.status != 'pending_payment':
logger.info(f"Order {order.id} status={order.status}, skip timeout")
return True
# 执行业务
close_order(order.id)
release_stock(order.id)
return True
5.4 监控与告警
生产环境必须对延迟队列建立监控体系,核心指标包括:
| 指标 | 获取方式 | 告警阈值 | ||
|---|---|---|---|---|
| 队列积压量 |
|
> 10000 | ||
| 最大延迟时间 |
最新任务 – 当前时间 |
> 60s(说明消费跟不上) | ||
| 死信队列大小 |
|
> 0 | ||
| Processing 集合大小 |
|
> 500 | ||
| 恢复任务次数 | 日志统计 | 持续增长说明消费者有问题 |
可以使用 Prometheus + Grafana 对接 Redis Exporter 实现可视化监控,配合 Alertmanager 设置告警规则。
5.5 Redis Cluster 兼容性
如果使用 Redis Cluster,Lua 脚本中操作的多个 Key 必须在同一个 slot 上,否则会报
1 | CROSSSLOT |
错误。解决方案是使用 Hash Tag——将队列相关的 Key 都加上
1 | {delay_queue} |
前缀:
1
2
3
4
5
6
7
8
9 # 使用 Hash Tag 确保 Key 在同一 slot
queue_key = "{delay_queue}:zset"
task_data_key = "{delay_queue}:data"
processing_key = "{delay_queue}:processing"
timeout_key = "{delay_queue}:processing:timeout"
dlq_key = "{delay_queue}:dlq"
# Redis Cluster 会根据 {} 内的内容计算 slot
# 所以上述 Key 都会落在同一个 slot 上
六、性能优化建议
6.1 批量处理
每次轮询取多条任务(batch_size=50~100),减少 Redis 交互次数。Lua 脚本本身就是批量操作,开销很小。
6.2 动态轮询间隔
消费者代码中已实现了动态间隔:如果队列中有任务但未到期,按距离最近任务的时间动态计算 sleep 时长,避免空轮询浪费 CPU 和 Redis 连接。
6.3 连接池配置
1
2
3
4
5
6
7
8
9
10 pool = redis.ConnectionPool(
host='127.0.0.1',
port=6379,
db=0,
max_connections=50,
socket_timeout=5,
socket_connect_timeout=3,
retry_on_timeout=True,
)
r = redis.Redis(connection_pool=pool)
6.4 数据分片
当单队列积压超过十万级时,考虑按任务类型分片到不同队列:
1
2
3
4
5
6 # 按业务类型分片
queue_order = RedisDelayQueue(r, queue_key="{delay}:order")
queue_retry = RedisDelayQueue(r, queue_key="{delay}:retry")
queue_notify = RedisDelayQueue(r, queue_key="{delay}:notify")
# 每个队列独立消费,互不影响
七、方案总结与最佳实践
最后,总结 Redis ZSet 延迟队列的生产级最佳实践清单:
- 用 Lua 脚本保证原子性:取任务+移除必须原子化,避免多消费者重复消费。
- 引入 Processing 中间态:任务取出后进入 processing 集合,ACK 后才彻底清理,防止消费者崩溃丢任务。
- 定时恢复超时任务:后台协程每 30 秒扫描 processing 超时的任务,重新入队。
- 业务层实现幂等:永远假设任务可能被重复执行,用 Redis SETNX 或数据库唯一约束做幂等。
- 死信队列兜底:超过最大重试次数的任务进入 DLQ,人工介入处理。
- 指数退避重试:失败重试时按 2^n 秒退避,避免下游服务被打爆。
- Hash Tag 兼容 Cluster:Redis Cluster 环境下务必用
1{tag}
前缀确保 Key 同 slot。
- 建立监控告警:队列积压量、死信队列、processing 超时数三个指标必须告警。
- 按业务分片:大流量场景按任务类型拆分到不同队列,避免互相阻塞。
Redis ZSet 延迟队列的精髓在于:用最简单的基础设施解决 80% 的延迟任务需求。它不需要引入额外的中间件,不需要复杂的配置,几十行 Lua 脚本就能实现生产级可靠性。当你的延迟任务量级达到百万级 TPS 时,再考虑迁移到 RocketMQ 或 Kafka 时间轮方案也不迟。
汤不热吧