上一篇 下一篇 分享链接 返回 返回顶部

在 Debian 的香港服务器上用 Redis Streams 扛稳在线答题实时统计:队列配置全攻略与高可用优化实战

发布人:Minchunlin 发布时间:2025-09-21 11:26 阅读量:754
那天凌晨 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%,就像我那晚做的——把每一次“抽搐”当成和系统对话的机会。祝你稳定、低延迟、也别熬太多夜。
目录结构
全文