背压、过载保护与故障恢复
每秒进入 900 个任务,系统只能完成 500 个,未完成任务就会持续增加。把这些任务移到内存队列、连接池等待区或消息 Broker,只改变了等待位置;只有减少新接收的工作、提高实际完成速度,或按业务规则终止一部分任务,积压才可能下降。
流量在哪一层变成等待
到达率、并发和队列是不同的量
外部请求
└── 准入:接受或拒绝
└── 等待队列:尚未开始的任务
└── 工作并发:已经开始、尚未结束
└── 下游连接/数据库/外部服务
└── 完成、失败或取消到达率用“每秒多少个请求”表示,并发用“同时有多少个未结束操作”表示,队列长度则数尚未开始的任务。相同到达率下,下游耗时从 20 毫秒变成 200 毫秒,会让更多操作同时占用线程、连接和内存。只限制 QPS,仍可能在下游变慢时耗尽并发资源。
入口的 offered、admitted、rejected 应分开计数。被拒绝的请求如果在客户端重试,会再次进入 offered;不能把尝试次数当作新增业务数量。completed 也要说明统计的是成功业务、终止任务还是包括拒绝的 HTTP 响应,计算容量时不能混用。
一个请求还可能依次等待多个队列:服务器分派线程、业务执行器、数据库连接池和数据库锁。把每处的等待上限都设成两秒,不会得到两秒的端到端响应;排队时间与实际工作时间会相加。
Little's Law 帮助核对平均量
对同一统计范围,在长期平均存在且系统可持续处理流量的条件下,Little's Law 给出:
平均系统内数量 L = 实际进入该范围的平均速率 λ × 平均停留时间 W
平均队列数量 Lq = 相应入队速率 λq × 平均排队时间 Wq若某执行阶段平均每秒处理 500 个任务、平均停留 0.2 秒,平均在途数量约为 100。只统计等待队列时,应使用排队时间,不能把包含业务执行的总延迟直接代入。公式及其平均量含义见 Little 原论文。
这个关系不能直接算出 p99,也不能证明一个正在无限积压的系统可以长期运行。任务大小、到达突发和处理时间波动都会改变尾部等待。用“目标等待时间 × 安全服务率”估算队列容量,可以得到一个初始候选值,还要用突发和慢任务负载验证;队列里每个任务保留的对象与请求体也计入内存预算。
四种控制方式作用不同
| 控制 | 决定什么 | 对调用方的影响 |
|---|---|---|
| 背压 | 下游当前允许上游继续发送多少数据 | 能配合的生产者减速或暂缓发送 |
| 速率限制 | 一段时间内接受多少请求或工作量 | 超出速率时等待、拒绝或降级 |
| 并发准入 | 同时允许多少个操作占用某种资源 | 无名额时尽早拒绝或进入有限等待 |
| 负载削减 | 容量不足时放弃哪些工作 | 按业务价值丢弃、采样、合并或返回错误 |
公网客户端未必配合减速,遥测设备也可能无法暂停采样。此时服务只能在接入处限制资源并选择超额数据的处理方式。付款指令和界面刷新不能采用同一种丢弃策略;最新温度可以覆盖旧温度,账务增量则需要可靠保存或明确拒绝。过载处理的容量与取舍见 Google SRE 负载管理。
用需求控制异步流中的元素
数据向下游流动,需求向上游传递
Reactive Streams 包含 Publisher、Subscriber、Subscription,以及同时承担发布和订阅的 Processor。订阅后,Subscriber 获得 Subscription,再通过 request(n) 累加需求:
Publisher ── onSubscribe / onNext / onComplete / onError ──► Subscriber
Publisher ◄──────────── request(n) / cancel ────────────── Subscriber
request(2) → 最多再发 2 个元素
request(3) → 再增加 3 个额度,累计最多发 5 个
cancel() → 请求结束该订阅关系并释放相关资源在一条订阅关系上,累计 onNext 数量不能超过累计请求数量;需求按元素计算,不按字节或线程计算。正需求由 request(n) 表达,零初始需求可以通过暂不调用 request 实现;直接调用 request(0) 是无效请求。Long.MAX_VALUE 可表示不再限制需求。Reactive Streams 1.0.4 规范
onComplete 和 onError 是终止信号,不占用元素额度;完成与失败不能接连作为两种终态发送。cancel 要求发布者最终停止信号和清理引用,但取消传播可能落后于正在途中的数据。取消本地订阅,也不能撤销已经提交的远端数据库事务。
背压协议不要求每个阶段都有独立线程。同步源可以在调用 request 的线程上发送元素,异步源则通过调度与队列交接。耗时阻塞工作若占住事件循环,需求数字再小也无法让其他连接获得运行机会;阻塞适配与调度器选择见 Reactor 线程与调度。
运行真实 request 与 cancel
Linux 宿主准备 Bash、unzip 和可用的 Docker Engine/Compose v2,使用已获准访问 Docker daemon 的非 root 用户。下载背压与准入实验包,保存为 distributed-backpressure-lab.zip,在新目录解压:
mkdir ds14-pressure-work
unzip distributed-backpressure-lab.zip -d ds14-pressure-work
cd ds14-pressure-work/distributed-backpressure
docker version
docker compose version
bash run-tests.sh 25Docker 应显示客户端和服务端。run-tests.sh 创建当前用户可写的 .m2 缓存,以该用户的 UID/GID 运行 maven:3.9.12-eclipse-temurin-25,显式挂载源码和缓存,不挂载 Docker socket 到 Maven 容器。编译目标为 Java 17,Reactor core/test 固定 3.8.7,JUnit 固定 6.0.3;bash run-tests.sh 17 使用 Temurin 17 运行同一组测试。Reactor 版本入口见 3.8.7 发行说明。
源码包含 6 个测试:3 个流需求测试、2 个执行器测试、1 个真实回环 HTTP 测试。正常套件报告 Tests run: 6, Failures: 0, Errors: 0, Skipped: 0。如果依赖不可达,使用企业批准的 Maven 仓库、镜像代理,或搬运已校验的同版本镜像和依赖缓存;不要关闭 TLS 或跳过负例。缓存不可写时修正本实验目录的属主,避免切换为 root 构建。
DemandTest 对一个实际的 Flux.range 订阅,初始需求为 0,然后请求 2 个、再请求 3 个,最后取消:
StepVerifier.create(source, 0)
.expectSubscription()
.then(() -> assertEquals(0, produced.get()))
.thenRequest(2).expectNext(1, 2)
.then(() -> assertEquals(2, produced.get()))
.thenRequest(3).expectNext(3, 4, 5)
.thenCancel()
.verify(Duration.ofSeconds(5));source 上的 doOnRequest、doOnNext 和 doOnCancel 记录实际信号,字段与完整测试位于包内 DemandTest.java。输出为 requested=5 produced=5 cancelled=true。未请求的后 5 个元素没有被发送;这是该同步源的受控结果,不能据此假定任何网络源都能立即取消。StepVerifier可以检查事件顺序、初始需求和限定等待时间。
buffer、prefetch 与 flatMap 会改变需求形态
buffer(4) 把 4 个源元素组合成一个 List。下游请求 2 个 List,上游需要提供 8 个元素,所以另一个测试得到 requested_groups=2 upstream_elements=8。如果只在最终订阅处看见 request(2),就把整个流的内存需求算成 2 个原始对象,会漏掉中间转换。
publishOn 等异步交接会使用预取队列;flatMap(mapper, concurrency, prefetch) 的 concurrency 限制同时订阅的内部序列数量,prefetch 调整向这些序列请求的额度。它们都不是每秒 QPS 上限。一个内部序列可以产生多个元素,也可能在订阅时开启远端请求;内存和下游连接预算必须把这些行为一起计算。concatMap 适合要求顺序且可以接受较低并发的处理。Reactor 需求变换、Flux API
onBackpressureBuffer(2) 采取另一种方式:向上游请求无限需求,再以本地有限缓冲承接。测试在初始需求为 0 时先让源填满缓冲,随后请求并取走前两个元素,观察溢出错误:
buffer_capacity=2 upstream_unbounded=true overflow=true这项有界重载在溢出后停止上游,并在已缓冲元素排空后交付错误;无参数的缓冲重载则可能积累无界数据。不同 drop/latest/error 策略会丢失不同信息,必须与业务允许的损失匹配。相应行为可在固定版本的 FluxOnBackpressureBuffer 实现中对应到请求、队列和终止分支。
流内的限额还可能被应用绕开。若 onNext 只把任务提交给另一个无界线程池便立即返回,上游看到的只是“提交任务很快”,真正的数据库工作仍在持续堆积。异步工作应作为流的一部分返回,或由目标执行器的容量明确拒绝;额外调用一次 fire-and-forget subscribe() 会使原订阅难以协调它的失败和取消。
线程池与 HTTP 入口分别限制什么
有界执行器怎样接受和拒绝任务
ThreadPoolExecutor 接收新任务时,先考虑核心线程,再尝试入队;队列不能接收时,才尝试扩展到最大线程数。线程也不能增加时,调用拒绝策略。线程池停止后同样可能拒绝提交,应区分停机与过载。Java ThreadPoolExecutor
实验采用以下参数,任务通过闩锁保持未完成:
new ThreadPoolExecutor(
1, 2, 0, TimeUnit.MILLISECONDS,
new ArrayBlockingQueue<>(1),
new ThreadPoolExecutor.AbortPolicy());| 提交次序 | 工作线程 | 排队任务 | 实际行为 |
|---|---|---|---|
| 第 1 个 | 1 | 0 | 创建核心工作线程 |
| 第 2 个 | 1 | 1 | 进入唯一队列位置 |
| 第 3 个 | 2 | 1 | 队列已满,创建额外线程 |
| 第 4 个 | 2 | 1 | 抛出 RejectedExecutionException |
PoolTest 等待前两个实际工作线程进入闩锁,再断言拒绝;闩锁让观察期间的工作线程与队列保持已知状态。随后释放闩锁,确认已接受的 3 个任务全部完成并关闭线程池。结果为 bounded workers=2 queued=1 rejected=1。
无界队列的对照将参数改成 core=1、maximum=4、new LinkedBlockingQueue<>()。第一个任务被闩锁挡住后,再提交 8 个任务,结果为 unbounded core=1 maximum=4 workers=1 queued=8。队列持续接受新任务,执行器就没有因队列满而增到 4 个线程。释放闩锁后,9 个任务由该线程完成,队列随之清空。
拒绝策略决定后续责任。AbortPolicy 让提交者处理失败;CallerRunsPolicy 让调用线程执行任务,可以拖慢某些同步生产者,却可能阻塞 HTTP 事件循环或持锁线程;DiscardPolicy 会丢弃任务。使用 Future 时还要处理被丢弃任务的完成状态,避免调用方一直等一个永远不会执行的任务。执行器选择与关闭细节见线程池、队列与拒绝。
业务名额耗尽时,返回明确的 HTTP 响应
线程池拒绝不一定会自动转换成 HTTP 503。请求可能在业务处理前就因连接或分派队列满而断开。需要稳定业务拒绝语义时,应在能够写响应的位置做准入,同时为接入层本身设置独立上限。
实验服务器用一个 Semaphore(1) 保护业务段,拿不到名额立即返回 503、Retry-After: 1 和 {"error":"BUSY"};已接受请求等待一个最长 30 秒的实验闩锁。这里只用等待控制名额占用,没有访问真实数据库。Semaphore
下面的 HTTP 操作需要 curl 7.76.0 或更新版本,以支持 --fail-with-body。完成前面的编译后,在同一解压目录启动:
docker compose up -d
docker compose logs --tail=20 app
docker inspect ds14-pressure-app --format '{{.Config.User}}'
curl -q --noproxy '*' --silent --show-error --fail-with-body --max-time 3 \
http://127.0.0.1:18440/state日志出现 listening=8080 business_slots=1 dispatch_threads=4,身份为 10001:10001,状态为 {"active":0,"admitted":0,"rejected":0,"completed":0}。刚启动时连接尚未就绪,可在日志确认监听后重试状态查询;其他错误先停止后续步骤。ClassNotFoundException 表明还未成功编译,回到 run-tests.sh;端口冲突则换独立实验环境,不关闭别人的服务。
运行镜像为 eclipse-temurin:25.0.4_7-jdk,应用文件只读挂载,只有 /tmp 可写;不另建镜像。宿主仅发布 127.0.0.1:18440,无 TLS,/release 是无认证的本机实验控制入口,不应对外暴露。JDK HttpServer负责实际网络连接;HTTP 分派配置为 4 个线程、8 个等待位置,与唯一业务名额分开。
下面几组命令应连续执行,避免在人工阅读期间超过 30 秒的等待上限。也可以执行 bash http-lab.sh 自动跑同样过程;已经运行过时,先 docker compose down 再重启全新状态。
先发出一个等待中的请求,并确认它占据名额:
mkdir -p http-output
curl -q --noproxy '*' --silent --show-error --fail-with-body --max-time 40 \
http://127.0.0.1:18440/work > http-output/first.json &
FIRST_PID=$!
for attempt in $(seq 1 60); do
state=$(curl -q --noproxy '*' --silent --show-error --fail-with-body --max-time 2 \
http://127.0.0.1:18440/state)
[[ "$state" == *'"active":1'* ]] && break
sleep 0.05
done
[[ "$state" == *'"active":1'* ]]只有状态判断成功后才继续。此时再发一个请求,分别判断传输是否成功、HTTP 状态和拒绝字段:
if code=$(curl -q --noproxy '*' --silent --show-error --max-time 5 \
-D http-output/rejected.headers -o http-output/rejected.json \
-w '%{http_code}' http://127.0.0.1:18440/work); then
test "$code" = 503
else
printf '%s\n' '传输失败,未得到预期的过载响应'
false
fi
tr -d '\r' < http-output/rejected.headers | grep -Ei '^Retry-After: 1$'
test "$(< http-output/rejected.json)" = '{"error":"BUSY"}'curl 传输应退出 0,状态为 503,响应包含指定字段与正文。这个负例有意不使用 --fail-with-body,避免把 HTTP 503 和网络连接失败混为一类;其他预期成功的请求仍保留该选项。-q 放在第一个选项以忽略 curlrc,--noproxy '*' 排除代理路径,选项说明见 curl 手册。
释放闩锁,再检查请求完成与后续恢复:
curl -q --noproxy '*' --silent --show-error --fail-with-body --max-time 5 \
--data '' http://127.0.0.1:18440/release
wait "$FIRST_PID"
test "$(< http-output/first.json)" = '{"result":"done"}'
curl -q --noproxy '*' --silent --show-error --fail-with-body --max-time 5 \
http://127.0.0.1:18440/work
curl -q --noproxy '*' --silent --show-error --fail-with-body --max-time 5 \
http://127.0.0.1:18440/state两次成功请求返回 {"result":"done"},最终状态为 {"active":0,"admitted":2,"rejected":1,"completed":2}。如果先出现 504/LAB_WAIT_EXPIRED,说明人工控制超过了实验时限,应关闭该实例并重启后再试;增加 curl 等待时间不能恢复已经结束的业务等待。
业务处理的成功、超时和异常都经过 finally 释放名额,随后分派线程发送响应;响应写失败也不会留下被占用的业务名额。completed 统计业务段完成,不能据此判定客户端已经收到响应。将名额释放只写在成功分支,少量异常就足以把服务永久变成“忙”。不过这个小服务没有监听客户端断连来立即中断等待,客户端主动退出后名额最多仍可能占用到闩锁释放或 30 秒超时;生产接入要将取消接到实际工作。
503 表达服务暂时无法处理,Retry-After 可以建议等待秒数;这并不预留下一次请求的名额。HTTP 语义规定了相应字段和状态。针对某个调用者的速率配额,429 Too Many Requests 通常更贴切,见 RFC 6585。
本实验的请求量远低于 HTTP 接入容量。若把它当压测服务持续洪泛,分派队列或 TCP 接入可能先饱和,客户端会断连而非每次收到 503。业务信号量保护的范围不能扩张成整个网络服务器的承诺。
切断重试与级联故障
一次超时可能增加下一轮负载
重试是新的资源消耗。若三层调用各允许最多 4 次尝试,一次原始请求在最坏情况下可能触发下游 64 次尝试。应选择知道错误语义、能保持原操作 ID 的一层负责重试,并将重试次数、等待和实际执行都计入原请求预算。Google SRE 级联故障
截止时间包含排队和执行。任务出队时已经过期,应结束等待并归还名额;不能再为它重新分配一份完整下游超时。取消 Future 通常只是请求取消,任务、驱动或远端服务还需要支持相应中断;结果未知的写操作沿原幂等键查询或重试,相关协议见超时、重试与幂等。
退避让失败请求延后,抖动把它们分散到不同时间。仍需限制全部客户端合计产生的额外尝试,例如单独分配重试令牌,令牌不足便停止重试或交给持久任务稍后处理。只为每个请求设置“最多三次”,在大面积失败时依旧可能放大接近三倍流量。
收到 Retry-After 后,先检查剩余预算能否覆盖等待与下一次执行,再按接口契约调度。不要让所有客户端恰好在同一秒醒来;服务端建议的等待是最低准备信号之一,不是保证成功的预约。
熔断、隔离与降级各有对象
熔断器根据失败或慢调用进入 OPEN,暂时阻止实际调用;HALF_OPEN 只放有限探测请求,恢复后再扩大。统计窗口、最小样本和探测并发要与流量规模匹配,不能把低流量的一次失败直接解释成系统性不可用。
并发隔离把不同依赖或业务类别分到独立名额与执行器,避免一个慢供应商用完所有连接。关键写入、普通查询和后台回填可以设置不同预算,但隔离过细会闲置资源,优先级队列也可能让低优先级任务永久饥饿。需要保留最低份额、任务年龄和人工升级条件。
降级必须改变实际工作量。返回缓存、减少推荐项或延迟统计可以少访问下游;在所有失败时同步调用另一个同容量服务,可能只是转移过载。高成本请求还应按任务重量计费:一项扫描十万行的查询与一次主键读取,各占一个请求额度并不公平。
跨进程后继续寻找等待位置
RabbitMQ 的 prefetch 限制未确认投递。消费者若收取后立即 ACK,再把任务放进无界本地队列,Broker 看不到真正尚未完成的工作。ACK 应与业务完成语义一致,预取数结合处理并发、消息大小和数据库能力选择;prefetch=0 表示不设该限制。RabbitMQ 预取说明
HTTP/TCP 的流量窗口主要约束字节传输。请求体已经被全部读进内存后,后端仍可能无限创建业务任务;限制连接数也可能把排队移到连接池的等待者。需要逐层核对在途数量、等待容量、等待期限和拒绝行为,不能只检查最外层有没有限流插件。
扩容应用实例也会扩大连接池总量。单实例 50 个数据库连接,10 个实例就是最多 500 个;如果数据库已经饱和,继续扩到 20 个实例可能使锁竞争与尾延迟更严重。虚拟线程降低部分线程成本,但不会增加数据库连接或外部接口额度。
定位积压并控制恢复速度
把现象对应到具体等待点
| 观察 | 常见解释 | 下一步 |
|---|---|---|
| 队列增长,工作线程都忙 | 到达超过完成,或任务变慢 | 查服务时间、线程栈与下游等待 |
| 队列增长,线程数停在 core | 队列一直可入,maximum 未被触发 | 查队列类型和执行器参数 |
| 应用 CPU 不高,连接等待很长 | 数据库连接或外部调用占满 | 查连接池等待者、持有时间、慢 SQL |
| original 稳定,attempts 激增 | 多层或同步重试放大 | 限制重试入口,核对退避与总预算 |
| 队列数量下降,最老任务不动 | 热点、优先级饥饿或单个坏任务 | 按对象、年龄与失败原因拆分 |
本机 HTTP 服务仍在时,可以查询状态和线程:
curl -q --noproxy '*' --silent --show-error --fail-with-body --max-time 3 \
http://127.0.0.1:18440/state
docker stats --no-stream ds14-pressure-app
docker exec ds14-pressure-app jcmd 1 Thread.print在首次请求尚未释放时,业务线程会停在 CountDownLatch.await;释放后,空闲分派线程可能停在 ArrayBlockingQueue.take。前者是任务内等待,后者是正常等下一项任务,仅凭 WAITING 状态无法判断卡死。容器 PID 1 是该 Java 进程,诊断使用与进程相同的 UID;真实服务应先确认目标 PID,避免套用这个编号。
这些检查分别提供业务在途、容器资源和线程位置。生产还需要队列长度、最老等待时间、成功服务率、拒绝原因、连接池等待以及分位延迟。状态计数不是完整监控系统,平均值也不能定位少数长期未完成的对象。
净排空速度扣除当前新工作
设积压为 120000 个任务、恢复后的安全完成能力为每秒 500 个。在任务成本相近、该能力能够持续且不重复计数的假设下:
| 外部尝试 / 秒 | 拒绝 / 秒 | 新接收 / 秒 | 完成 / 秒 | 净排空 / 秒 | 理想排空时间 |
|---|---|---|---|---|---|
| 900 | 400 | 500 | 500 | 0 | 不会排空 |
| 900 | 600 | 300 | 500 | 200 | 600 秒 |
| 900 | 600 | 300 | 400 | 100 | 1200 秒 |
净排空 = 实际完成 − 新接收。拒绝掉 400 并不会凭空获得每秒 200 个的恢复余量;500 个新接收正好用完 500 个完成能力。客户端重试如果重新进入队列,也必须计入新接收量。任务重量差异很大时,改用估计工作量或分别计算各类任务,不能直接相减条数。
积压中可能已经有过期请求。浏览器早已关闭的查询可以结束,能够从源数据重建的派生索引可以选择重建,而已接受的付款任务必须保持可查询结果或进入人工处理。删除队列里的消息之前,要先确定终止它的业务含义。
恢复坡度由下游反馈决定
先恢复依赖与少量真实业务请求,确认返回内容、写入结果和必要缓存已经可用;再逐档扩大新请求与积压消费。每档观察尾延迟、拒绝率、下游等待、实际完成率和最老积压年龄,覆盖足以看到新负载影响的时间窗口。
当完成率不再增长、下游等待持续上升或最老积压失去下降趋势时,暂停增加并发,必要时回退上一档。恢复控制同时约束新流量、重试和后台回填,避免三个来源各自认为还有全部容量。健康探测应便宜且有独立预算;把过载直接当成存活失败持续重启,会进一步减少工作能力。
应用升级或回退也要带上这些设置:队列容量、超时、并发、重试次数和默认预取变化都会改变系统负载。检查旧版本能否接手已经接受的持久任务,先停止领取新任务,再限时等待或交还未完成工作;不要只看进程成功退出。
关闭实验实例
所有 HTTP 观察完成后,在解压目录执行:
docker compose down只回收 ds14-pressure 应用容器和专用网络。实验工作与计数在内存中,关闭后不保留;源码、测试报告和 Maven 缓存仍在当前目录。无需删除其他项目容器、卷或全局镜像缓存。
