Learning
VOL. XIII · NO. 24 · Elixir · 01 JAN 1970

Redix bounded LPUSH / 背压 / 队列监控

Elixir 编程 · 01 JAN 1970 · 7 min read · 1,367 words
· · ·

Redis 的 LIST 是最简单的队列,但生产环境必须解决三个问题:连接池、上游背压、运行时监控。Redix 是 Elixir 生态首选的 Redis 客户端。

本课目标
学完后,你应该能用 Redix 启动连接池;用 bounded LPUSH 实现有上限的入队(背压);用 XLEN / LRANGE + Telemetry 监控队列深度;区分 BLPOP vs BRPOP vs 轮询的取舍。

一、Redix 是什么

Redix 是 Elixir 生态最常用的 Redis 客户端。它走 Erlang 的二进制协议(RESP),不是文本协议,效率高。和 Finch 类似,它是一个 supervisor + connection pool,但只有一个 Redis 连接池(不像 Finch 那样按 host 拆分)。

二、启动 Redix 实例

# config/config.exs
config :my_app, :redis,
  url: System.get_env("REDIS_URL", "redis://127.0.0.1:6379/0"),
  pool: [size: 10, max_overflow: 5]

# application.ex
children = [
  {Redix, Application.get_env(:my_app, :redis) ++ [name: :redix]}
]

关键参数:

  • size:常驻连接数。
  • max_overflow:突发时可创建的额外连接(用完即归还)。
  • name:注册名,调用时通过它拿连接。

三、基本命令

# 同步命令
Redix.command!(:redix, ["PING"])
#=> "PONG"

# 异步管道
Redix.pipeline!(:redix, [
  ["SET", "user:1", Jason.encode!(%{name: "alice"})],
  ["EXPIRE", "user:1", 3600]
])
#=> ["OK", 1]

# 异步:返回 {:ok, pid}
{:ok, pid} = Redix.command(:redix, ["GET", "user:1"])
同步 vs 异步
  • command!/2pipeline!/2:阻塞当前进程直到 Redis 返回。
  • command/2pipeline/2:返回 {:ok, pid},命令结果通过消息发给 pid。
  • 在 Phoenix 请求处理里优先用异步,避免阻塞 Cowboy/Bandit acceptor。

四、LPUSH / RPOP:经典队列

Redis 的 LIST 可以当队列:

  • LPUSH key value:从左侧入队。
  • RPOP key:从右侧出队(FIFO)。
  • BLPOP key timeout:阻塞出队,空时等待。
  • BRPOP key timeout:阻塞出队,反向(RPUSH + BLPOP = FIFO)。
# 生产者
def enqueue(queue, payload) do
  Redix.command!(:redix, ["LPUSH", queue, payload])
end

# 消费者(同步)
def pop(queue) do
  case Redix.command!(:redix, ["RPOP", queue]) do
    nil -> :empty
    v -> {:ok, v}
  end
end

这个版本的问题是没有背压——消费者慢、生产者快,列表会无限增长,Redis 内存爆掉。

五、bounded LPUSH:让入队有上限

Redis 5+ 提供 LMPOP / LPOP 的 COUNT 参数,但入队没有原生上限。要做”满了就拒”,用 Lua 脚本原子地检查长度:

# 启动时把脚本加载到 Redis
SCRIPT = """
local n = redis.call('LLEN', KEYS[1])
if n >= tonumber(ARGV[1]) then
  return -1
end
redis.call('LPUSH', KEYS[1], ARGV[2])
return n + 1
"""

{:ok, sha} = Redix.command!(:redix, ["SCRIPT", "LOAD", SCRIPT])
#=> {:ok, "a4b1..."}

然后在生产者端:

def bounded_push(queue, payload, max) do
  case Redix.command!(:redix, ["EVALSHA", @sha, "1", queue, max, payload]) do
    -1 -> {:error, :queue_full}
    n  -> {:ok, n}
  end
rescue
  Redix.Error ->  # NOSCRIPT:脚本被淘汰,重新加载
    {:ok, new_sha} = Redix.command!(:redix, ["SCRIPT", "LOAD", SCRIPT])
    :persistent_term.put({__MODULE__, :sha}, new_sha)
    bounded_push(queue, payload, max)
end
bounded LPUSH 的意义
  • 队列长度有了上限,Redis 内存不会爆。
  • 返回 :queue_full 时上层可以降级(同步处理、丢弃、写磁盘缓冲)。
  • 这是最便宜的"背压"——比维护 broker 简单得多。

六、消费者:阻塞 vs 轮询

消费者有两种选择:

策略优点缺点
BLPOP低延迟,空时占用少占一个连接不放
轮询连接可复用CPU 抖动 / 延迟取决于轮询间隔

如果一个 BEAM 节点只跑几个 worker,用 BLPOP 最简单:

def blocking_consumer(queue) do
  case Redix.command!(:redix, ["BRPOP", queue, "5"]) do
    nil -> :timeout
    [_key, value] -> handle(value)
  end
end

如果跑很多 worker,每个 worker 占一个 Redis 连接代价太大,用 XREAD + Stream 更合适;或者干脆走 Oban。

七、队列监控:XLEN / Telemetry

队列深度是必须监控的指标。Redis 提供 XLEN(LIST 长度)直接读:

def depth(queue) do
  Redix.command!(:redix, ["LLEN", queue])
end

# 周期性上报
def monitor_loop(queue, interval \\ 5_000) do
  Process.send_after(self(), :tick, interval)

  receive do
    :tick ->
      d = depth(queue)
      :telemetry.execute([:my_app, :queue, :depth], %{value: d}, %{queue: queue})
      monitor_loop(queue, interval)
  end
end

[:my_app, :queue, :depth] 接到 Prometheus reporter,就能看到 P95 队列深度。配合告警阈值:

defp alert?(depth) when depth > 10_000, do: {:alert, "queue too deep"}
defp alert?(_), do: :ok
监控要点
  • 深度(XLEN):队列堆积了多少任务。
  • 入队速率(每秒 LPUSH 数)。
  • 出队速率(每秒 RPOP 数)。
  • 消费延迟(任务从入队到出队的时间)。
  • 错误率(payload 反序列化失败、handler 抛错)。

八、消费者错误处理

消费者处理失败时,不能简单 RPOP 后丢弃:

  • DLQ:失败时 RPOPLPUSH source dlqBRPOPLPUSH source dlq timeout
  • 重试:失败时 LPUSH source payload 把任务放回队首。
  • 指数退避:失败时把任务推到 queue:retry:1s / queue:retry:5s 等等。
def safe_consume(queue, dlq) do
  case Redix.command!(:redix, ["BRPOPLPUSH", queue, dlq, "5"]) do
    nil -> :timeout
    value ->
      try do
        handle(value)
        Redix.command!(:redix, ["LREM", dlq, "1", value])  # 处理成功,移出 DLQ
      rescue
        e ->
          Logger.error("consumer failed: \#{inspect(e)}")
          # 让任务留在 DLQ,等后续 retry 工具拉回主队列
      end
  end
end
Oban 的取舍
如果任务需要持久化、重试、唯一性、监控面板,直接用 Oban。自己维护 bounded queue + DLQ 是 Oban 不适用的场景:流式、临时、跨语言消费者。

九、Redix 的连接复用与背压

如果入队 QPS 高,单连接会成为瓶颈。配置更大的 pool:

config :my_app, :redis,
  url: "redis://...",
  pool: [size: 50, max_overflow: 20]

但 pool 大不代表吞吐无限。Redis 自身的 QPS(单实例通常 5-10 万)才是天花板。再大就要做:

  • Redis Cluster 分片。
  • 多个独立 Redis 实例(按业务切)。
  • 把队列降级到本地(ETS + bounded push),跨节点同步用 gossip 或 pull。

十、生产清单

  • 队列深度有上限(bounded LPUSH),永远不要无限 LIST。
  • 每个 Redis pool 配监控:连接数、命令延迟、错误率。
  • 消费者失败必须留痕(DLQ + 日志),不能丢。
  • 生产 / 消费两端都要限速(Hammer / 自实现令牌桶)。
  • Redis 重启或切换时,DLQ 里的任务要可重新入队。
  • 不要在 BRPOP 里放敏感超时——它会占连接,连接池会被它吃光。

十一、测验

1

Redix 默认的连接模式是?

2

LPUSH + RPOP 实现的是?

3

bounded LPUSH 用什么实现?

4

BRPOPLPUSH source dlq timeout 的作用是?

5

BLPOP 的主要代价是?

6

Redis 单实例的 QPS 上限大约是?

7

Redix.command/2command!/2 的区别?

8

队列深度持续增长但消费速率不变,说明?

9

失败任务应该?

10

如果任务需要持久化、重试、唯一约束,最适合的方案是?

**下一步:**Redis 队列适合临时、流式场景;Oban 适合持久化重试唯一性。下一课看 Oban 的生命周期细节:cron / queue / retry / discarded / 失败通知。