一个响应最多输出 32 MiB,每次只产生 16 KiB。客户端收到响应头以后暂停读取,服务端生产计数停在 49 块;客户端随后只读一块就关闭真实 socket,服务端流进入终态,关闭计数为一。

这个实验测到了当前链路中的需求反馈与取消传播。49 块已经远大于客户端实际读取的一块,差额可以停留在流、服务器和内核缓冲中。把应用 Source 的 inputBuffer 设成一,不能把整条网络链路理解成“只允许提前产生一个元素”。

生产速率、消费速率与待发送字节

流式响应首先改变的是生产方式。有限列表先把所有对象构造完再交给 Source,虽然输出按元素传输,列表占用的内存已经发生。实验的 Source.range 只产生整数索引,map 在获得下游需求时分配一个 16384 字节数组,上限为 2048 块,不构造完整字节列表。

先忽略对象开销,以字节计算未消费数据。假设生产速率为每秒 8 MiB、消费速率为每秒 1 MiB,持续两秒且生产端没有减速,待发送数据会增加 14 MiB。这是速率差的积分,不是已经测出的 JVM 堆占用。如果可用缓冲只有 512 KiB,差额不能永久累积,系统需要降低生产速率、拒绝新增数据或终止连接。

1
2
3
4
生产: p bytes/s ──> [应用缓冲] ──> [服务器缓冲] ──> [内核发送缓冲] ──> 消费: c bytes/s
<──── 需求与可写性逐段反馈 ────
待消费增量近似值 = max(0, p - c) × 持续时间
应用元素数量上限 ≠ 全链路字节上限 ≠ 进程 RSS

背压描述生产方如何响应下游需求,容量描述允许多少未消费数据停留在系统中。若一个 map 接收到需求后发起无界数量的后台任务,或者先把全部查询结果放进集合,外层流遵守需求协议仍不足以约束内部资源。

本篇使用 Play 3.0.6、Pekko Streams 1.0.3、Scala 2.13.15 和 JDK 21。app/controllers/StreamLabController.java 提供数据端点,app/streamlab/StreamLedger.java 记录生产与终止。账本最多保留 128 条记录,容量不足时只回收已终止记录;全部处于活动状态时拒绝继续分配,避免监控实验自身形成无界 Map。

HTTP 流的需求、缓冲与终止

从 Source 到 HttpEntity

Play 固定提交为 2e56aff7d4e7a74af61e4bd39ec9e3ed7f300cd6。HttpEntity.Streamed 同时持有字节 Source、可选内容长度和可选内容类型。这三个值有不同用途:Source 决定如何取得字节,长度帮助服务器划定消息边界,类型描述内容。

构造 Result 时通常还没有消费 Source。流被服务器接到输出端并物化后,才产生这一次运行的状态。把计数器绑定在共享的静态 Source 外部,再跨请求重复使用,会混淆多次物化;实验为每个请求建立独立 Entry,并使用唯一 id 查询它。

Pekko 固定提交为 4f77c8108aaf548a65531d2c8807da13dbba8146。Java DSL Source.range 构造范围流;同文件的 watchTermination 将终止 Stage 交给物化值组合函数。回调拿到的是流运行的完成信号,不是客户端业务代码已经处理完全部字节的确认。

下面的完整类可放入累计工程 app 目录编译。它只展示有限生产与物化完成接口;带 id、异常注入和账本的实际端点保存在 StreamLabController。示例以 JDK 21 验证。

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
import java.util.Optional;
import java.util.concurrent.CompletionStage;
import java.util.function.Consumer;
import org.apache.pekko.Done;
import org.apache.pekko.stream.Attributes;
import org.apache.pekko.stream.javadsl.Source;
import org.apache.pekko.util.ByteString;
import play.http.HttpEntity;

public final class FiniteExport {
public static HttpEntity create(
int chunks, boolean known,
Consumer<CompletionStage<Done>> onMaterialized) {
if (chunks < 1 || chunks > 2048) {
throw new IllegalArgumentException("chunks must be 1..2048");
}
Source<ByteString, ?> data = Source.range(0, chunks - 1)
.map(index -> {
byte[] block = new byte[16384];
java.util.Arrays.fill(block, (byte) ('a' + index % 26));
return ByteString.fromArray(block);
})
.withAttributes(Attributes.inputBuffer(1, 1))
.watchTermination((materialized, done) -> {
onMaterialized.accept(done);
return materialized;
});
return new HttpEntity.Streamed(
data, known ? Optional.of(chunks * 16384L) : Optional.empty(),
Optional.of("application/octet-stream"));
}
}

ByteString 的构造还可能涉及复制,因此“每块 16 KiB”限定的是负载大小,不能精确等同于每次分配量。即使只保留少量可达数组,GC 尚未回收的对象也可能影响堆占用。要分析峰值,需要同时观察存活对象、分配速率、服务器缓冲与进程级指标。

已知长度与未知长度的实际响应头

真实客户端使用 HTTP/1.1,请求 64 块,每次 read 16384 字节,读取后等待 2 毫秒。两种响应均接收 1048576 字节,账本均为 producedChunks=64、terminated=true、closeCount=1。

构造方式 Content-Length Transfer-Encoding Connection 本次消息结束依据
Streamed,长度已知 1048576 无 未显式 close 指定字节数
Streamed,长度未知 无 无 close 响应连接关闭
第19篇显式 chunked SSE 无 chunked 按服务器连接策略 chunked 编码结束

因此,当前后端并没有把每一个未知长度的 Streamed 自动改写成 chunked。ok().chunked(...) 与 new HttpEntity.Streamed(..., Optional.empty(), ...) 是不同的实体表达。文章中的结论来自本次响应头,不能扩展到 HTTP/2、Netty 或其他版本。

长度写错同样会破坏协议边界。声明的字节数大于实际输出时,客户端可能等待剩余内容或报告截断;小于输出时,额外字节无法按这个响应的原声明解释。流式生成能够提前确定总长度时,应从同一输入约束计算长度,避免使用近似值。

暂停读取与实际取消

取消场景使用默认 2048 块,即 33554432 字节上限。客户端先取得响应头,暂停 600 毫秒,通过另一条 HTTP 连接查询账本,然后读取 16384 字节并关闭响应与连接。第二条连接继续有限轮询终态,避免关闭后失去所有观察通道。

场景 暂停时已生产块数 终态块数 终态已生产字节 closeCount
已知长度,关闭 socket 49 49 802816 1
未知长度,关闭 socket 49 49 802816 1

49 是这次机器、缓冲和调度下的采样值。测试断言只要求暂停时尚未产生全部 2048 块、关闭后最终终止,以及 closeCount=1,不把 49 固化为跨平台常量。返回 Result、Source 已物化、生产一块、客户端读到一块,必须作为四个时点分别记录。

两种取消的 failure 字段在本次记录中为空。下游取消可以让 watchTermination 的完成 Stage 正常完成,不能把空 failure 解释为完整传输成功。判断完整传输仍需比较预期长度、实际生产量和客户端接收量。

进程 RSS 在同轮全场景执行中采样 237 次,基线为 227648 KiB,峰值为 240544 KiB。采样覆盖的还有 SSE、WebSocket 与文件实验,不能把二者差值直接归因为这个 Source。它只能排除这次有限窗口中已经采样到的更高值,不能证明恒定内存或所有瞬时峰值。

异常、停止与普通 Future 的差异

异常分支在索引 512 处抛出 IllegalStateException,前面已经生产 8388608 字节。响应头已经是 HTTP 200,随后读取出现不完整消息;服务器不能在同一响应体中补发另一个 HTTP 500 来改写已经发送的状态。

这类接口需要定义业务数据是否允许部分结果。例如 CSV 导出可以附带独立校验文件或导出任务状态,但不能仅凭最初的 200 宣称整个导出完成。若消费者需要原子结果,可以先生成有限文件再提供下载,这会改变延迟、磁盘容量与清理责任。

独立 JUnit 场景物化 SSE 后停止应用。账本记录 producedChunks=1、terminated=true、closeCount=1,异常为 AbruptStageTerminationException。这个实验证明应用停止会终止被测运行时中的活动流;它没有打开数据库游标,不能据此声称数据库连接也已释放。资源必须在拥有它的 IO 阶段设置关闭路径,第21篇另外验证文件 IO。

普通异步 action 则有不同结果。/streamlab/work 安排 750 毫秒后的一个合成工作单元,返回原始 CompletableFuture。客户端确认 workStage=scheduled 后关闭 TCP;后续轮询观察到 workStage=completed、producedChunks=1、futureCancelled=false。

1
2
普通 action: scheduled ── 客户端关闭 TCP ── 工作执行 ── 原 Future completed
流式 body: materialized ─ 客户端关闭 TCP ── 下游取消 ── watchTermination 完成

网络连接断开没有自动给普通任务增加协作取消协议。真实任务若包含副作用,需要独立传递截止时间、取消意图和业务幂等标识;流终止信号也只能清理已经纳入该运行生命周期的资源。

重跑与改动练习

共享累计工程已另行验收:39项JUnit与stage通过,8个流/文件class与生产jar字节一致;DEV/PROD各198次既有HTTP和五组真实socket通过。共享记录在evidence/batch18-21/shared-socket/observations.json;隔离样本保留原测量值,不将两轮耗时、RSS、块数混为同一轮。

从第00篇取得累计工程,并按该篇设置 JDK 21 与 sbt launcher。从 play-lab 目录执行:

1
2
bash sbtw clean test stage
python3 lab/stream_checks.py --evidence evidence/batch18-21/replay

脚本自行启动 staged server,占用实验端口 19018,并在 finally 停止服务器。每个请求 id 唯一,流账本最多保留 128 项;重跑应使用新的证据目录。暂停、关闭、失败和 Future 对照的数据位于 observations.json 的 httpStreams。完整历史网络观测保存在 evidence/batch18-21/isolated/run-final-io/observations.json。

本轮隔离累计工程的 clean test 为 36 项 JUnit 全通过,其中 5 项属于 StreamLabTest;stage 通过。它是隔离实验收据,不代替后续共享工程的测试数量或页面验收。HTTP/2、生产代理、数据库流式事务和长期负载均未在这组证据中运行。

反例题:在构造 Source 前执行 List<byte[]> all = ...,再设置 inputBuffer(1, 1),能否沿用本篇的内存结论?不能;列表已经把全部负载变成可达对象,背压无法撤销之前的分配。

改动练习:保持 32 MiB 总量,把每块大小分别设为 4 KiB 和 64 KiB,再固定客户端每秒读取量。同时记录暂停时元素数与字节数、RSS 样本、关闭后的终态。若只比较元素数,块大小变化会把同一缓冲字节量误报成背压改善。

第二个练习是在错误注入前后分别读取响应头和固定长度的响应体。为“完整接收”“源失败”“客户端主动取消”建立三个独立字段,禁止用一个 closeCount 或一个 HTTP 状态码合并它们。

上一篇:缓存加载与失效。下一篇:SSE与慢客户端。