那天凌晨 00:42,香港机房 12U 的那台小钢炮又把我从 Slack 里“抓”了出来:实时答题统计的榜单卡死。直播间 1.8 万人同时刷题,后端 API 200 都是绿的,MySQL 也没报警,唯独榜单每隔十几秒“抽搐”一下。
我戴上工牌,坐进冷飕飕的机房过道,在机器嗡鸣声里开始了这次 Redis 队列的“拆解与重构”。这篇文章就是那天夜里到天亮,我一步步做过、踩过、修过的所有细节。
场景与目标
业务:在线答题系统(秒级实时榜单、题目正确率、用户连击数等)。
痛点:高并发时榜单抖动、统计延迟、偶发重复计数。
目标:
队列有序、可恢复,至少一次投递语义;
统计聚合**< 1 秒**延迟稳定;
可横向扩展消费者;
故障(如重启、Failover)后不丢、不重、可自愈。
1. 机房与硬件基线(真实参数)
| 组件 |
型号/规格 |
说明 |
| 服务器 |
Intel Xeon E-2288G 8C/16T @ 3.7GHz |
单路、睿频高,适合低延迟队列 |
| 内存 |
64GB DDR4 ECC |
Redis + Page Cache 充足 |
| 磁盘 |
2×1.92TB NVMe(RAID1,mdadm,ext4) |
AOF 落盘稳性 + 读写延迟低 |
| 网卡 |
1GbE 独享 |
业务与监控共用,VLAN 隔离 |
| OS |
Debian 12 (bookworm),内核 6.1 |
官方 redis-server 7.x |
| 时钟 |
chrony(对时 NTP 香港节点) |
统计窗口对齐,避免时间漂移 |
我把 Redis 与业务 API 放在同一机柜内的两台物理机上,通过机房内网互访,控制平面(Sentinel、监控)另起两台轻量实例做三节点仲裁。
2. 系统级调优(先把地基打平)
2.1 sysctl(网络 & FD)
/etc/sysctl.d/99-quiz-redis.conf:
net.core.somaxconn = 65535
net.ipv4.tcp_max_syn_backlog = 4096
net.ipv4.tcp_syncookies = 1
net.ipv4.tcp_tw_reuse = 1
net.core.netdev_max_backlog = 16384
net.ipv4.tcp_fin_timeout = 15
fs.file-max = 2000000
vm.overcommit_memory = 1
vm.swappiness = 1
应用:sysctl --system
2.2 ulimit 与 THP
/etc/security/limits.d/redis.conf:
redis soft nofile 1000000
redis hard nofile 1000000
禁用透明大页(THP)并持久化:
echo never | tee /sys/kernel/mm/transparent_hugepage/enabled
echo 'echo never > /sys/kernel/mm/transparent_hugepage/enabled' >> /etc/rc.local
chmod +x /etc/rc.local
2.3 NVMe 与文件系统
ext4 + noatime,nodiratime,discard(确认 NVMe 固件支持在线 trim)
RAID1(mdadm)避免单盘故障导致 AOF 脏写不可恢复
3. 安装与进程管理
apt update && apt install -y redis-server
systemctl enable redis-server
Debian 12 默认 redis-server 7.x,带 Redis Streams 与 ACL/TLS 能力,正合我意。
4. Redis 配置:为“队列 + 实时聚合”而生
/etc/redis/redis.conf(关键摘选与解释):
# 监听与安全
bind 127.0.0.1 10.0.12.10 # 机房内网 IP
protected-mode yes
port 6379
# 如果走公网或跨机柜,开启 TLS(见下节)
# tls-port 6380
# 持久化:AOF 为主,RDB 作为 AOF 前导
appendonly yes
appendfilename "appendonly.aof"
appendfsync everysec # 性能与持久化折中
no-appendfsync-on-rewrite yes # 避免重写期 fsync 抖动
auto-aof-rewrite-percentage 50
auto-aof-rewrite-min-size 256mb
aof-use-rdb-preamble yes
# 内存与过期策略(答题统计是热数据)
maxmemory 28gb
maxmemory-policy allkeys-lfu # 热点友好,避免冷键撑爆
hash-max-ziplist-entries 1024
hash-max-ziplist-value 256
# 队列语义:用 Streams + 消费组
stream-node-max-bytes 4096
stream-node-max-entries 100
# 客户端与 I/O
maxclients 200000
timeout 0
tcp-keepalive 60
client-output-buffer-limit normal 0 0 0
io-threads 4 # 7.x 开始支持(主要用于网络 I/O)
# 复制与故障切换
replica-read-only yes
repl-backlog-size 512mb
为什么选 Streams?
自带消费组(XGROUP),至少一次投递,支持重分配未确认消息(XAUTOCLAIM)。
自然支持多消费者并行且有序,比 LPUSH/BRPOP 更适合这类统计流水。
5. TLS 与 ACL(可选但我强烈建议)
5.1 生成自签证书(或接入内网 CA)
openssl req -x509 -nodes -newkey rsa:4096 -sha256 -days 3650 \
-keyout /etc/ssl/private/redis.key \
-out /etc/ssl/certs/redis.crt -subj "/CN=redis.internal"
cat /etc/ssl/certs/redis.crt /etc/ssl/private/redis.key \
> /etc/ssl/private/redis.pem
chmod 600 /etc/ssl/private/redis.*
redis.conf 追加:
tls-port 6380
port 0
tls-cert-file /etc/ssl/private/redis.pem
tls-key-file /etc/ssl/private/redis.key
tls-ca-cert-file /etc/ssl/certs/redis.crt
5.2 ACL:应用与聚合分权
/etc/redis/users.acl:
user default off
user quizapp on >p@ss-QuizApp ~stream:answers:* +xadd +xinfo +ping
user statsworker on >p@ss-Worker ~stream:answers:* +xreadgroup +xack +xautoclaim +xpending +xclaim +hset +hincrby +eval +ping
redis.conf 引用:
aclfile /etc/redis/users.acl
6. 队列模型设计(关键键名与结构)
提交流:stream:answers:{quizId}
消息体字段(例):
uid:用户 ID
qid:题目 ID
ok:是否正确(0/1)
ts:答题时间戳(ms)
sig:幂等签名(hash(uid, qid, ts))
消费组:group:stats
消费者命名:consumer-<hostname>-<pid>
聚合哈希(热统计在 Redis):
hash:quiz:{quizId}:stats
字段:total, correct, q:<qid>:total, q:<qid>:correct, u:<uid>:streak
持久化落库(每 500ms 批量 flush 到 MySQL/PostgreSQL):
避免每条消息触发 DB 写,降低抖动
7. 生产写入(API 侧,FastAPI 示例)
# app_producer.py
import time
import hashlib
import json
import redis
from fastapi import FastAPI
r = redis.Redis(host='10.0.12.10', port=6379, username='quizapp', password='p@ss-QuizApp')
app = FastAPI()
STREAM_PREFIX = "stream:answers:"
def sig(uid, qid, ts):
return hashlib.sha256(f"{uid}:{qid}:{ts}".encode()).hexdigest()[:16]
@app.post("/submit")
def submit(quiz_id: str, uid: str, qid: str, ok: int):
ts = int(time.time() * 1000)
payload = {
"uid": uid, "qid": qid, "ok": ok, "ts": ts, "sig": sig(uid, qid, ts)
}
sid = r.xadd(f"{STREAM_PREFIX}{quiz_id}", payload, maxlen=1000000, approximate=True)
return {"stream_id": sid, "queued": True}
MAXLEN ~ 1000000 控制流长度,避免无限增长。真实值按峰值流量与 AOF 容量评估。
8. 消费与统计(aioredis + Streams 消费组)
# worker_stats.py
import asyncio, time, os
import aioredis
from collections import defaultdict
STREAM = "stream:answers:{quiz}"
GROUP = "group:stats"
CONSUMER = f"consumer-{os.uname().nodename}-{os.getpid()}"
BATCH = 200
ACK_AFTER = 0.5 # 秒,聚合窗口
async def ensure_group(r, quiz_id):
try:
await r.xgroup_create(STREAM.format(quiz=quiz_id), GROUP, id="$", mkstream=True)
except aioredis.exceptions.ResponseError as e:
if "BUSYGROUP" not in str(e): raise
async def bump(r, quiz_id, bucket):
# 原子聚合:pipeline + HINCRBY;必要时用 Lua 做一次多键原子(示例下节)
key = f"hash:quiz:{quiz_id}:stats"
pipe = r.pipeline()
for field, inc in bucket.items():
pipe.hincrby(key, field, inc)
await pipe.execute()
LUA_DEDUP = """
-- 幂等:用 SETNX 保证 sig 唯一消费,过期防膨胀
local k = KEYS[1]; local sig = ARGV[1]; local ttl = tonumber(ARGV[2])
if redis.call('SETNX', k .. sig, 1) == 1 then
redis.call('PEXPIRE', k .. sig, ttl)
return 1
else
return 0
end
"""
async def consume_quiz(quiz_id):
r = await aioredis.from_url("redis://10.0.12.10:6379", username="statsworker", password="p@ss-Worker")
await ensure_group(r, quiz_id)
dedup = r.register_script(LUA_DEDUP)
buffer, last = defaultdict(int), time.time()
while True:
# 先处理 PEL(pending)中的超时消息
# XAUTOCLAIM Redis 6.2+:把 idle>5s 的消息拉回本消费者
try:
_, claimed = await r.xautoclaim(STREAM.format(quiz=quiz_id), GROUP, CONSUMER, min_idle_time=5000, start_id="0-0", count=BATCH)
for sid, fields in claimed:
# 幂等判断
ok = await dedup(keys=[f"dedup:quiz:{quiz_id}:"], args=[fields[b"sig"].decode(), 60000])
if ok == 1:
qid = fields[b"qid"].decode(); u = fields[b"uid"].decode()
isok = int(fields[b"ok"].decode())
buffer["total"] += 1
buffer[f"q:{qid}:total"] += 1
if isok:
buffer["correct"] += 1
buffer[f"q:{qid}:correct"] += 1
buffer[f"u:{u}:streak"] += 1
else:
buffer[f"u:{u}:streak"] = 0
await r.xack(STREAM.format(quiz=quiz_id), GROUP, sid)
except Exception:
pass
# 常规读取新消息
resp = await r.xreadgroup(GROUP, CONSUMER, streams={STREAM.format(quiz=quiz_id): ">"}, count=BATCH, latest_ids=None, timeout=200)
if resp:
for _stream, messages in resp:
for sid, fields in messages:
ok = await dedup(keys=[f"dedup:quiz:{quiz_id}:"], args=[fields[b"sig"].decode(), 60000])
if ok == 1:
qid = fields[b"qid"].decode(); u = fields[b"uid"].decode()
isok = int(fields[b"ok"].decode())
buffer["total"] += 1
buffer[f"q:{qid}:total"] += 1
if isok:
buffer["correct"] += 1
buffer[f"q:{qid}:correct"] += 1
buffer[f"u:{u}:streak"] += 1
else:
buffer[f"u:{u}:streak"] = 0
await r.xack(STREAM.format(quiz=quiz_id), GROUP, sid)
# 周期性 flush 到 Redis 哈希(再由另一个任务批量落库)
if time.time() - last > ACK_AFTER and buffer:
await bump(r, quiz_id, buffer)
buffer.clear()
last = time.time()
if __name__ == "__main__":
asyncio.run(consume_quiz("spring_festival_2025"))
关键点:
XAUTOCLAIM + PEL 兜底,至少一次;
Lua + SETNX 做幂等;
批量 HINCRBY,避免每条消息一把 DB 锁。
9. Lua 原子聚合(可替换 pipeline)
当统计维度多、跨键多时,推荐把幂等 + 聚合写在一个 Lua 里,一次 EVAL 吃下(示例思路):
-- KEYS[1] = dedup prefix, KEYS[2] = stats key
-- ARGV = sig, ttl, qid, uid, ok
local dpx = KEYS[1]; local stats = KEYS[2]
local sig = ARGV[1]; local ttl = tonumber(ARGV[2])
local qid = ARGV[3]; local uid = ARGV[4]; local ok = tonumber(ARGV[5])
if redis.call('SETNX', dpx .. sig, 1) == 0 then
return 0
end
redis.call('PEXPIRE', dpx .. sig, ttl)
redis.call('HINCRBY', stats, 'total', 1)
redis.call('HINCRBY', stats, 'q:'..qid..':total', 1)
if ok == 1 then
redis.call('HINCRBY', stats, 'correct', 1)
redis.call('HINCRBY', stats, 'q:'..qid..':correct', 1)
redis.call('HINCRBY', stats, 'u:'..uid..':streak', 1)
else
redis.call('HSET', stats, 'u:'..uid..':streak', 0)
end
return 1
10. 后台落库(MySQL/PostgreSQL 批量)
每 500ms 读取 hash:quiz:{quizId}:stats 的增量分片(或用 SCARD 追踪脏字段),拼批量 SQL 写入;
数据库表按 quiz_id + time_bucket(1s) 分区/索引,避免热点写。
11. 进程守护(systemd)
/etc/systemd/system/quiz-stats-worker.service:
[Unit]
Description=Quiz Stats Worker
After=network.target redis-server.service
[Service]
User=www-data
Environment=PYTHONUNBUFFERED=1
ExecStart=/usr/bin/python3 /srv/quiz/worker_stats.py
Restart=always
RestartSec=1
LimitNOFILE=1000000
[Install]
WantedBy=multi-user.target
12. Sentinel + 复制(高可用)
12.1 拓扑
主:物理机 A(Redis)
从:轻量实例 B、C(同城不同机柜)
Sentinel:A/B/C 各 1 份(3 节点仲裁)
12.2 从库配置(摘)
replicaof 10.0.12.10 6379
AOF 同样开启,防止主故障期间从提升后数据不可追。
12.3 sentinel.conf(每台)
port 26379
sentinel monitor quiz-redis 10.0.12.10 6379 2
sentinel down-after-milliseconds quiz-redis 5000
sentinel failover-timeout quiz-redis 60000
sentinel parallel-syncs quiz-redis 1
生产上我把 API 端的 Redis 连接改为通过 Sentinel 获取主节点地址(或使用内置的哨兵客户端),这样 Failover 不需要发版。
13. 监控与告警(我在夜里最依赖的灯)
redis-exporter + Prometheus + Grafana
关键监控项与阈值(经验值):
| 指标 |
阈值/告警 |
说明 |
instantaneous_ops_per_sec |
与基线相比 +50% 或 -50% |
负载异常 |
latency(ms) |
p99 > 2ms 持续 1 分钟 |
本地内网应很低 |
aof_current_size 增长率 |
每分钟 > 300MB |
检查流量与重写 |
sync_partial_err |
> 0 |
复制异常 |
blocked_clients |
> 100 |
Lua/慢查询或 IO 堵塞 |
mem_fragmentation_ratio |
> 1.8 |
内存碎片与内核态压力 |
rejected_connections |
> 0 |
maxclients 触顶 |
x_pending(每流) |
> 10000 且持续 |
消费积压,扩容 Worker |
14. 压测基准(机房当晚我做过的)
命令(示例):
redis-benchmark -h 10.0.12.10 -p 6379 -n 200000 -c 200 -d 64 -t XADD,XREAD
实测(摘录)(仅供量级参考,不同环境差异大):
| 项目 |
QPS(均值) |
p99 延迟 |
| XADD 64B |
~180k/s |
1.2ms |
| XREADGROUP 64B |
~150k/s |
1.6ms |
真正瓶颈常在消费者逻辑(聚合/落库)而非 Redis 本身。
我在线上把 BATCH = 200 调到 400,p99 延迟显著下降,但榜单刷新更平滑——两端择一平衡。
15. 常见坑 & 我是怎么救火的
AOF 重写抖动榜单
症状:榜单每隔几分钟卡一下。
原因:AOF 重写期间 IO 峰值,appendfsync 遭遇抖动。
解法:no-appendfsync-on-rewrite yes + NVMe RAID1;重写阈值调到 50%/256MB,错峰在业务低谷触发(后台定时 BGREWRITEAOF)。
Pending 队列越积越多
症状:XPENDING 显示大量 idle>10s。
原因:消费者实例重启、死亡未 ACK。
解法:周期 XAUTOCLAIM + 启动时先扫 PEL;增加 consumer 并行数。
重复计数
症状:排行榜比 DB 大 0.3% 左右。
原因:重投递/网络抖动,至少一次语义带来的重复。
解法:Lua + SETNX 幂等签 + TTL 60s;对关键榜单(前 1,000 名)再与 DB 做秒级回写校准。
maxmemory 打爆淘汰了热键
原因:默认 LRU 不友好。
解法:allkeys-lfu,热度模型更稳;分离冷门统计到 DB,仅在 Redis 保留窗口内(例如最近 5 分钟)热数据。
慢日志惊现 EVAL 超时
原因:Lua 做了过多字符串拼接、循环。
解法:Lua 里仅做幂等与少量 HINCRBY,复杂逻辑回到应用侧批处理。
时钟漂移导致窗口错位
解法:统一 chrony,对齐到香港 NTP 池,应用层所有窗口基于 Redis TIME。
16. 参数总表(我线上用了这些)
| 类别 |
参数 |
值 |
备注 |
| 持久化 |
appendonly |
yes |
AOF |
| |
appendfsync |
everysec |
折中 |
| |
no-appendfsync-on-rewrite |
yes |
降抖动 |
| |
aof-use-rdb-preamble |
yes |
快速恢复 |
| 内存 |
maxmemory |
28gb |
物理 64G 留足系统 |
| |
maxmemory-policy |
allkeys-lfu |
热点友好 |
| I/O |
io-threads |
4 |
网络 I/O |
| 流 |
stream-node-max-bytes |
4096 |
节点拆分 |
| |
stream-node-max-entries |
100 |
节点大小 |
| 客户端 |
maxclients |
200000 |
够用 |
| 安全 |
aclfile |
users.acl |
分权 |
| HA |
repl-backlog-size |
512mb |
复制稳态 |
| 网络 |
tcp-keepalive |
60 |
长连保活 |
17. 运维手册(我放在 README 里的常用命令)
# 观察流长度与消费
redis-cli XINFO STREAM stream:answers:xxx
redis-cli XPENDING stream:answers:xxx group:stats
# 手动重分配 pending
redis-cli XAUTOCLAIM stream:answers:xxx group:stats consumer-1 5000 0-0 COUNT 100
# AOF 体检
redis-cli INFO persistence | egrep 'aof_current_size|aof_last_bgrewrite_status'
# 慢日志
redis-cli CONFIG SET slowlog-log-slower-than 1000
redis-cli SLOWLOG GET 10
18. 验收标准(我给团队立的线)
峰值 2 万 QPS 写入,3 个 statsworker 实例,p99 < 2s(端到端从提交到榜单刷新);
Failover(Sentinel 模拟主挂)期间消息不丢、统计不回退;
24h 内 x_pending 峰值 < 1 万 且 1 分钟内回落。
天快亮的时候,直播的老师还在讲最后一套题。我把耳机从服务器上拔下来,看了一眼 Grafana:Ingress 平稳、XPENDING 贴地飞行、榜单线如拉尺。
第二天,产品同事问我:“你昨晚是不是动了什么魔法?排行榜像钉子一样稳。”
我笑了笑,其实哪有什么魔法,不过是把 Redis 当成可靠的消息中间件 + 热内存数据库去认真打磨:正确的队列模型、谨慎的持久化策略、克制的 Lua、以及对抖动的敬畏。
机房的风依旧冷,但这次是我先离开的。下次高峰来时,它会自己稳住。
20. 可复用的清单(拿走即用)
Debian 12 安装 redis-server 7.x
应用 sysctl、ulimit、禁用 THP
NVMe RAID1 + ext4 noatime
redis.conf:AOF(everysec)、LFU、IO 线程
采用 Streams + 消费组 + XAUTOCLAIM
Lua 幂等(SETNX + TTL)
批量 HINCRBY 聚合,500ms 落库
Sentinel 三节点 + 从库
redis-exporter 指标告警
压测与回归基线记录
如果你也在香港或其他节点扛实时榜单,照着这些步骤落地,70% 的坑你会天然避开。剩下的 30%,就像我那晚做的——把每一次“抽搐”当成和系统对话的机会。祝你稳定、低延迟、也别熬太多夜。