两份字节相同的 Controller、路由和配置,都能正确解析 JSON、传输数据和完成 WebSocket 回声。向正在处理请求的生产进程发送 SIGTERM 后,结果发生了分叉:Pekko HTTP 返回完整的 200 completed,Netty 客户端收到 RemoteDisconnected。两边的业务回调和 stop hook 却都执行了。

Play 的应用层抽象使多数业务代码可以跨后端复用,连接关闭、协议错误和终止阶段仍由具体传输实现参与决定。本章固定版本,用相同输入、真实 TCP 连接和进程信号比较这些边界。

固定构件,确认真正运行的后端

下载两套服务器后端工程与实验记录,核对SHA256SUMS。解压后的 play-electives/e04-backends 包含Pekko与Netty应用;接入后重跑4项JUnit、两套stage与36项网络观察,20个应用class与生产jar字节一致,见 evidence/e04/integration。WebSocket协议错误与在途停机差异均保留。

实验固定 Play 3.0.6、Scala 2.13.15、sbt 1.10.7、JDK 21.0.11。Pekko 工程启用 PlayJava,Netty 工程使用下面的构建声明:

1
lazy val root = (project in file(".")).enablePlugins(PlayJava, PlayNettyServer).disablePlugins(PlayPekkoHttpServer)

这是官方提供的后端选择方式。两种服务构件都进入 classpath 时,不能只凭项目名称判断选择结果,应显式选择 provider 并检查运行时配置。Netty 后端文档

本次两个 stage 目录分别只包含所选 Play 服务构件。实际 JAR 清单如下,完整文件名和 SHA-256 保存在 run-final/artifacts.json:

层次 Pekko 工程 Netty 工程
Play 服务适配器 play-pekko-http-server 3.0.6 play-netty-server 3.0.6
HTTP 实现 pekko-http-core 1.0.1 netty-codec-http 4.1.115.Final
Netty 流适配 不包含 netty-reactive-streams-http 3.0.3
运行时 provider PekkoHttpServerProvider NettyServerProvider
请求 Server 属性 absent netty

健康接口同时读回 play.server.provider、实际请求属性和 Controller 类的代码来源。生产日志记录监听地址;类来源指向 stage 应用 JAR;编译目录中的 Controller 及内部类与应用 JAR 条目逐字节相同。配置、构件和实际请求三条证据结合,才能排除“改了构建文件却仍在访问旧进程”的情况。请求属性缺失本身不能普遍识别 Pekko;这里还检查了 provider 和互斥构件。

两套工程仍然都有 Pekko Streams。选择 Netty 只替换 HTTP 服务后端,不会把应用返回的 Source 改成另一种流 API。

flowchart LR
  C["HTTP / WebSocket 客户端"] --> P["Pekko HTTP 1.0.1"]
  C --> N["Netty 4.1.115"]
  P --> PA["Play Pekko 服务适配器"]
  N --> NA["Reactive Streams HTTP + Play Netty 适配器"]
  PA --> A["相同路由 / Filter / BodyParser / Controller"]
  NA --> A
  A --> S["Result / Source / WebSocket Flow"]
  S --> PA
  S --> NA

这个分层解释了为什么 JSON 解析结果可以一致,而传输错误仍可能不同。BodyParser 接收到应用层字节流后执行共同逻辑;在它前面的 HTTP 解码、连接关闭,在它后面的响应编码和发送,都有后端参与。

同一套输入验证到哪一层

两个工程的 BackendController.java、routes、application.conf 和 JUnit 文件逐字节相同。每边执行 2 个 JUnit,无跳过;随后分别从 target/universal/stage/bin 启动独立生产 JVM。JUnit 验证路由、JSON 和有限响应;后端差异由真实网络验证,不能用 route(app, request) 替代。

JSON parser 内存上限设为 1 KiB。合法对象返回 200,语法错误返回 400,错误 Content-Type 返回 415,超过上限的 JSON 返回 413,两边一致。这些输入已经通过各自 HTTP 解码器,随后进入相同 Play parser。没有覆盖请求行畸形、冲突的 Content-Length 或 HTTP request smuggling,不能从四个状态码推出整个协议解析器等价。

响应流每块为 16 KiB,间隔约 30 ms,完整场景总计 4 块。区分三种响应形式:

应用响应形式 两边实际 HTTP/1.1 framing 完整正文
Streamed,已知长度 Content-Length: 65536 65536 个零字节
Streamed,未知长度 Connection: close,无长度字段 相同
显式 chunked Transfer-Encoding: chunked 相同

三类正文 SHA-256 均为 de2f256064a0af797747c2b97505dc0b9f3df0de4f489eac731c23ae9ca9cc31。未知长度不等于显式 chunked;HttpEntity.Streamed 与 chunked(source) 的选择会影响 HTTP/1.1 消息结束方式。测试对字段名转为小写后断言,因为 HTTP 字段名不区分大小写。早期驱动将 Netty 的小写字段误判为缺失,失败记录保留在 run-01。

取消场景改为最多 128 块,客户端读完首块后关闭响应和 TCP 连接。随后通过独立健康连接轮询服务端状态,要求流终止、真实 FileChannel 已关闭、对应文件已删除,且生产量小于 128。最终计数为:

场景 Pekko 生产块数 Netty 生产块数
已知长度取消 3 3
未知长度取消 3 3
显式 chunked 取消 2 3
SSE 读到第 3 个 id 后取消 6 3

这些计数包含异步缓冲和调度影响,不能当作固定预取容量。可以确认的是:连接取消传播到了源,资源在应用停止前已释放。stop hook 是兜底清理,因此不能等进程退出后才检查资源,否则会掩盖流取消清理失效。

失败场景在已发送响应头后令 Source 抛出 IllegalStateException。两边 HTTP 状态已经是 200,客户端读取正文失败,服务端终态记录异常并关闭文件。错误发生得太晚时,不能重新发送一个 500 响应来替换已经发出的 200;调用方需要校验正文是否完整,不能只看响应头。

可编译的观测 Controller

下面是两套工程实际编译的完整 Java 文件。资源文件是实验用的生命周期探针,不保存业务数据;生产代码不应直接暴露这些诊断接口。实验为每个 id 限定格式、禁止复用,并把条目总数限制为 32,避免探针自身无限增长。

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
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
package controllers;
import com.typesafe.config.Config;
import jakarta.inject.Inject;
import jakarta.inject.Singleton;
import java.nio.channels.FileChannel;
import java.nio.file.*;
import java.time.Duration;
import java.util.Optional;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;
import org.apache.pekko.actor.ActorSystem;
import org.apache.pekko.stream.javadsl.Flow;
import org.apache.pekko.stream.javadsl.Source;
import org.apache.pekko.util.ByteString;
import play.api.mvc.request.RequestAttrKey;
import play.http.HttpEntity;
import play.inject.ApplicationLifecycle;
import play.libs.F;
import play.libs.Json;
import play.mvc.*;
@Singleton
public final class BackendController extends Controller {
private final Config config;
private final ActorSystem system;
private final ConcurrentHashMap<String, Entry> entries = new ConcurrentHashMap<>();
private final Path directory;
@Inject public BackendController(Config config, ActorSystem system, ApplicationLifecycle lifecycle) {
this.config=config; this.system=system;
try { directory=Files.createDirectories(Path.of(config.getString("lab.workdir"))); }
catch (Exception error) { throw new IllegalStateException(error); }
lifecycle.addStopHook(() -> {
for (Entry e: entries.values()) e.close(null);
event("stop-hook");
return CompletableFuture.completedFuture(null);
});
}
private void event(String value) {
try { Files.writeString(directory.resolve("journal"), value+"\n", StandardOpenOption.CREATE, StandardOpenOption.APPEND); }
catch (Exception error) { throw new IllegalStateException(error); }
}
private final class Entry {
final AtomicInteger produced=new AtomicInteger();
final FileChannel file;
final Path path;
volatile boolean closed;
volatile String failure="";
Entry(String id) throws Exception {
path=directory.resolve(id+".resource");
file=FileChannel.open(path,StandardOpenOption.CREATE_NEW,StandardOpenOption.WRITE);
}
synchronized void close(Throwable error) {
if (closed) return;
if (error!=null) failure=error.getClass().getSimpleName();
try { file.close(); Files.deleteIfExists(path); closed=true; }
catch (Exception cleanup) { throw new IllegalStateException(cleanup); }
}
}
private synchronized Entry open(String id) {
if (!id.matches("[a-z0-9_]{1,32}") || entries.containsKey(id) || entries.size()>=32) throw new IllegalArgumentException("finite unique id required");
try { Entry e=new Entry(id); entries.put(id,e); return e; }
catch (Exception error) { throw new IllegalStateException(error); }
}
public Result health(Http.Request request) {
return ok(Json.newObject().put("provider",config.getString("play.server.provider"))
.put("serverAttribute",request.attrs().getOptional(RequestAttrKey.Server().asJava()).orElse("absent"))
.put("classOrigin",getClass().getProtectionDomain().getCodeSource().getLocation().toString()));
}
@BodyParser.Of(BodyParser.Json.class)
public Result json(Http.Request request) { return ok(request.body().asJson()); }
public Result state(String id) {
Entry e=entries.get(id); if(e==null) return notFound();
return ok(Json.newObject().put("produced",e.produced.get()).put("closed",e.closed).put("fileOpen",e.file.isOpen()).put("fileExists",Files.exists(e.path)).put("failure",e.failure));
}
public Result data(String id, int chunks, boolean known, int failAt, boolean chunked) {
if (chunks<1 || chunks>128) return badRequest();
Entry e=open(id);
Source<ByteString,?> source=Source.range(0,chunks-1).throttle(1,Duration.ofMillis(30)).map(i -> {
if(i==failAt) throw new IllegalStateException("synthetic-stream-failure");
e.produced.incrementAndGet(); return ByteString.fromArray(new byte[16384]);
}).watchTermination((mat,done)->{done.whenComplete((v,x)->e.close(x));return mat;});
if (chunked) return ok().chunked(source).as("application/octet-stream");
return ok().sendEntity(new HttpEntity.Streamed(source,known?Optional.of(chunks*16384L):Optional.empty(),Optional.of("application/octet-stream")));
}
public Result sse(String id) {
Entry e=open(id);
Source<ByteString,?> source=Source.range(1,40).throttle(1,Duration.ofMillis(40)).map(i -> {
e.produced.incrementAndGet(); return ByteString.fromString("id: "+i+"\ndata: value-"+i+"\n\n");
}).watchTermination((mat,done)->{done.whenComplete((v,x)->e.close(x));return mat;});
return ok().chunked(source).as("text/event-stream");
}
public WebSocket ws(String id) {
return WebSocket.Text.acceptOrResult(request -> {
if(!"http://localhost".equals(request.header("Origin").orElse(""))) return CompletableFuture.completedFuture(F.Either.Left(forbidden()));
Entry e=open(id);
Flow<String,String,?> flow=Flow.of(String.class).map(text->{e.produced.incrementAndGet();return "echo:"+text;})
.completionTimeout(Duration.ofSeconds(8)).watchTermination((mat,done)->{done.whenComplete((v,x)->e.close(x));return mat;});
return CompletableFuture.completedFuture(F.Either.Right(flow));
});
}
public CompletionStage<Result> work() {
CompletableFuture<Result> result=new CompletableFuture<>();
event("work-start");
system.scheduler().scheduleOnce(Duration.ofMillis(800),()->{event("work-complete");result.complete(ok("completed"));},system.dispatcher());
return result;
}
}

watchTermination 观察流的终态,Entry.close 串行、幂等地关闭文件。清理操作若失败,状态不会被标记为成功;轮询也不会误报文件已经消失。完成 Future 的成功或失败与 HTTP 客户端完整收包是不同观测面,所以驱动还检查客户端正文、异常和实际关闭帧。

SSE 在这里是带事件格式的 HTTP 流。客户端检查 text/event-stream 和递增的 id,再主动取消。实验未实现 Last-Event-ID 重放,也没有经过反向代理,因此没有验证浏览器重连、代理缓冲和断点续传语义。

WebSocket 的正常关闭与协议错误

两个生产服务都先检查 Origin。不符合限定值的握手返回 403,状态接口证明没有创建资源。合法握手返回 101,驱动校验 Sec-WebSocket-Accept,发送带掩码文本并收到 echo:hello,再发送 Close 1000。两边均返回 Close 1000,文件清理成功。

反例发送了未掩码的客户端数据帧。该输入违反 WebSocket 客户端帧规则,故不能用普通回声成功来预测结果:

观测 Pekko Netty
客户端收到的终态 Close 1002 TCP EOF,没有 Close 帧
服务端流失败记录 空 CorruptedWebSocketFrameException
文件关闭、删除 成功 成功

Netty 原始日志包含:

1
WebSocket communication problem: received a frame that is not masked as expected

这证明异常输入确实到达协议处理链。Play 的共同 WebSocketFlowHandler 在输入流失败时记录通信问题并传播失败;底层解码器将协议问题转换为什么终止事件,会影响线上可见结果。固定源码:WebSocketFlowHandler

本次没有进一步定位 Netty 解码器内部每个关闭分支,也没有抓包证明中间所有帧。可以确定驱动没有收到 Close 帧便遇到 EOF,不能把它记录成 Close 1002 或 1006;1006 是本地用于描述异常断开的状态值,不能伪造成线上接收的关闭码。最初假设两边都返回 1002 的失败保留在 run-02,修正观察口径后的全量复跑保留在 run-final。

SIGTERM:业务完成不保证响应送达

/work 收到请求后记录 work-start,800 ms 后由 ActorSystem 调度器记录 work-complete 并完成响应 Future。驱动确认开始事件后立即发送 SIGTERM,同时保留客户端连接。两边使用相同 5 秒服务终止预算与 7 秒 service-requests-done 阶段预算。

sequenceDiagram
  participant C as 客户端
  participant S as 服务后端
  participant A as 应用回调
  participant L as Lifecycle
  C->>S: GET /work
  S->>A: 记录 work-start,安排 800 ms 回调
  Note over S: JVM 收到 SIGTERM,停止监听
  alt Pekko HTTP
    S->>S: binding.terminate 等待现有请求
    A-->>S: completed Result
    S-->>C: 200 completed
  else Netty
    S--xC: 关闭 non-server channels
    A-->>S: completed Result,客户端已断开
    S->>S: 关闭 event loop
  end
  S->>L: 应用 stop hook
  L->>L: 记录 stop-hook

最后一轮的原始观测摘录:

1
2
3
4
pekko: status=200 body=completed clientError=null exit=143
netty: status=null body="" clientError=RemoteDisconnected exit=143
both: work-start -> work-complete -> stop-hook
both: remainingResourceFiles=0 listenerClosed=true

这次从发送信号到退出约为 Pekko 0.920 秒、Netty 2.992 秒。退出耗时包含关闭调度与 quiet period,不能解读为吞吐或延迟排名。Netty 退出更晚也没有保住这条 HTTP 响应。

固定源码给出了差异依据。Pekko 服务先调用 binding.unbind(),随后在 service-requests-done 阶段调用带预算的 binding.terminate(...)。PekkoHttpServer 3.0.6

Netty 服务先关闭监听 channel,然后在 service-requests-done 阶段关闭剩余非监听 channel,再执行 event loop 的 shutdownGracefully。这里的 graceful shutdown 包含事件循环的关闭协议,不能直接等同于“等待所有 HTTP 响应发送完毕”。NettyServer 3.0.6

上线策略因此需要把摘流、等待和进程信号分开设计。先让负载均衡停止分配新请求,再观察应用在途计数并等待预算,最后发送 SIGTERM。此次回调完成不能证明数据库事务提交;客户端失败也不能证明业务没有完成。若请求会产生不可重复副作用,仍需幂等键或结果查询协议来解决响应丢失后的重试问题。

复跑与证据范围

完整工程在 examples/play-electives/e04-backends,原始证据在 examples/play-electives/evidence/e04/isolated。从项目根目录进入示例,设置自己的 JDK 21 路径,按 README 下载并核验 sbt launcher 后执行:

1
2
3
4
cd examples/play-electives/e04-backends
export PLAY_LAB_SBT_LAUNCHER="$PWD/.tools/sbt-launch-1.10.7.jar"
python3 build_checks.py --evidence evidence/build
python3 http_checks.py --evidence evidence/run

构建脚本执行 clean、test、stage;网络脚本使用临时端口和只属于本次运行的资源目录,结束时停止两个 JVM。最终结果是 4 个 JUnit 执行、0 跳过,36 条网络与进程观测通过。正文完整 Java 文件还单独使用 JDK 21、stage classpath 编译。历史失败与最终成功分目录保存,避免把驱动修复前后的日志混成一次实验。

已运行的是本机 HTTP/1.1、有限正文、真实 TCP 取消、两种 WebSocket 关闭路径和 SIGTERM。TLS、HTTP/2、反向代理、Linux epoll 实际运行、多节点、长时间并发、吞吐与内存基准均为 NOT_RUN。虽然 Netty 包中包含 Linux native 构件,本次 macOS 运行不能证明加载或使用了它们。

两个练习

延迟摘流练习。 在现有工程增加显式 readiness 和在途计数,先摘流,再等待 /work 完成,最后发送 SIGTERM。分别断言客户端结果、业务事件、停止 hook 和端口关闭;再将等待预算缩短到 100 ms,保存超时路径,比较是否仍出现响应丢失。

协议错误练习。 增加超长控制帧或非法 UTF-8 文本,保留握手原文、发送帧十六进制和接收到的 Close 帧或 EOF。分别记录服务端流终态与资源关闭;不要把客户端库计算的本地关闭状态当作线上帧中的状态码。这些新增输入尚未执行。

上一篇扩展:E03:Ebean 与 Hibernate 集成。下一篇扩展:E05:Pekko Actor。生产摘流与强杀后的恢复见第 30 篇。