客户端收到 id 为 101、102、103 的三个事件后关闭连接,再携带 Last-Event-ID: 103 建立新请求,得到 104、105、106。这次连续性来自服务器保存的事件集合与游标选择规则;仅给事件加一个递增 id,无法提供相同的重放能力。

同一个有限事件源经过两种本地代理模式时,首个事件的交付时间也不同。立即转发模式几乎立刻交付,按两秒条件刷出的缓冲模式约等待 2.279 秒。控制器已经生成事件,不代表下游每一层都已把事件交给客户端。

一条事件可以跨越多个读取操作

Server-Sent Events 使用文本格式表达服务端事件。应用需要同时处理 HTTP 传输边界和 SSE 事件边界。网络上的一个 read 可能只得到半行,也可能包含多个事件;使用“一次 read 对应一个事件”的解析方法会把 TCP 分段误当业务分段。

本篇的每个 sample 事件包含两行 data、一行 id 和一行 event,空行结束事件。首次消息格式如下:

1
2
3
4
5
event: sample
id: 101
data: line-one-101
data: line-two-101

两个 data 行属于同一事件。若业务原始字符串包含换行,直接拼接 data: 与原始字符串,会让第二行脱离 data 字段。应使用能逐行编码的 EventSource.Event,而不是自行假设内容永远为单行。

实验客户端以 readline 累积字段,到空行才提交事件。它保留 data 数组,因此能明确断言每个事件含两个 data 字段。实际业务处理器还需要决定如何把这些行还原为文本,不能把字段扫描结果与最终业务对象混为一谈。

SSE 编码、保留记录与客户端游标

EventSource.flow 编码的范围

Java 入口位于 Play 3.0.6 固定提交 2e56aff7d4e7a74af61e4bd39ec9e3ed7f300cd6 的 core/play-java/src/main/java/play/libs/EventSource.java。flow() 把 Event 的 formatted 结果转换为 ByteString,Event.formatted 再委托 Scala 编码实现。

core/play/src/main/scala/play/api/libs/EventSource.scala 负责字段和逐行 data 格式。Java Flow 的存在没有增加持久化、确认、重试或重放;这些都不在字符串编码函数的职责范围内。

Pekko Streams 固定为 1.0.3,提交 4f77c8108aaf548a65531d2c8807da13dbba8146。累计工程中的 SseLabController 使用有限 Source、throttle 与 takeWithin,在物化时注册终止回调。对应 Java DSL 实现可查 Source.scala。

下面的完整类可放入累计工程 app 目录编译,返回一个两行 data 的有限事件流。它展示编码和评论心跳的顺序,不承担重连游标逻辑;包含固定保留集合的完整端点见 app/controllers/SseLabController.java。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
import java.time.Duration;
import org.apache.pekko.stream.javadsl.Source;
import org.apache.pekko.util.ByteString;
import play.libs.EventSource;
import play.mvc.Result;
import play.mvc.Results;

public final class FiniteSse {
public static Result events() {
Source<ByteString, ?> body = Source.range(101, 120)
.throttle(1, Duration.ofMillis(450))
.map(id -> EventSource.Event
.event("line-one-" + id + "\nline-two-" + id)
.withId(Integer.toString(id))
.withName("sample"))
.via(EventSource.flow())
.intersperse(ByteString.fromString(": heartbeat\n\n"))
.takeWithin(Duration.ofSeconds(10));
return Results.ok().chunked(body)
.as("text/event-stream")
.withHeader("Cache-Control", "no-cache");
}
}

这里的心跳是两个有限事件之间插入的 SSE comment,不带 id,也不推进客户端业务游标。它能够验证客户端跳过评论行,但没有演示“长时间完全没有业务事件时仍定时发送”的独立保活服务。生产接口若存在长时间空闲,心跳生成与超时策略需要另外设计和测试。

throttle 限制的是源侧生成节奏。代理仍可以先消费并保存多条事件,再一次交付。客户端看到突发的事件批次时,应同时检查中间层缓冲,不能直接断言 throttle 没有生效。

明确有限重放的游标规则

实验只保留固定的 20 个事件,id 为 101 到 120。它们是可重复读取的内存对象,不是连接每生成一次就丢弃的临时编号。首次请求不带游标时按 100 处理,取所有 id 大于游标的事件。

Last-Event-ID 响应 解释
缺省或 100 从 101 开始 固定日志的起点
103 从 104 开始 已收到 103,重放剩余记录
99 410 游标早于允许起点,无法承诺所缺记录仍在
120 204 这份有限日志已全部读完
121 或无法解析的文本 400 超出当前日志或格式错误

过期游标不能悄悄重置为最新位置,否则客户端会把跳过的数据误认为已经连续收到。实际系统可以选择返回明确过期错误、要求客户端获取全量快照,或提供单独的补偿读取入口;每一种都需要相应客户端处理。

游标也不是授权信息。多用户事件流应先确定调用者可见的事件范围,再解释该范围内的游标。同一个整数出现在两个租户中,不代表它们可以共享重放位置。本实验使用公开的合成数据,未实现租户隔离与持久化日志。

1
2
3
4
5
保留集合: 101 102 103 104 105 106 ... 120
连接 A: └────── 接收 ───┘ TCP close
保存游标: 103
连接 B: Last-Event-ID=103 ──> 104 105 106 ...
过期输入:99 ──> 410,不悄悄跳到最新事件

将游标持久化放在业务处理之前,会出现“游标已推进、业务未处理”的丢失窗口。放在业务处理之后,则可能在断电重连时重复处理最后一条。需要恰好一次业务效果时,消费者必须把去重或业务幂等纳入设计,SSE 传输本身无法替代这一步。

三事件断线重连的真实观测

Python 标准库客户端建立 HTTP 连接,逐行接收三个完整事件,然后实际关闭连接对象和响应。服务端终态通过另一个有限轮询请求读取。随后客户端新建连接,显式发送 Last-Event-ID: 103,而不是复用原连接或直接调用 controller 方法。

观察项 结果
首次三条 id 101、102、103
每条 data 行数 2
首次读取期间评论行 至少 2
新连接三条 id 104、105、106
两个连接的终止记录 terminated=true,closeCount=1
过期、读完、未来、非法游标 410、204、400、400

这份结果只覆盖固定集合仍在内存中的重连。没有运行跨进程重启、事件淘汰竞争、分区切换与消费者崩溃恢复,不能写成一般性的“断连不丢消息”。已经发送的 id 提供定位能力,能够重放的记录和明确的过期协议才提供恢复能力。

两秒缓冲怎样改变首事件时间

本地代理有两种模式。direct 每读到一块就转发并 flush;buffered 将块加入 bytearray,在达到容量阈值或距上次 flush 超过两秒时刷出。每次上游读取最多 1024 字节,阈值预留这一块的空间,缓冲断言不超过 65536 字节。

这个两秒检查在读到上游新数据时触发,不是独立定时器。因此在当前每 450 毫秒产生事件的输入下,实际等待会大于两秒;如果上游完全静默,代码不会凭空唤醒一次读操作。本地示例的有限输入和 socket 超时保证它不会无限挂起。

模式 首事件等待,秒 完整接收事件数
direct 0.00001825 20
buffered 2.27906479 20

计时从客户端取得代理响应对象后、开始解析 SSE 行时起算。首事件可能已经进入 socket 缓冲,所以 direct 的微秒级数值不能解释成完整的端到端 HTTP 延迟。测试断言比较两种模式至少一秒的差距,不要求复现表内小数。

两个模式都收到了全部 20 个事件,却有不同的交付延迟。只核对最终事件数量会漏掉实时性退化;只测首事件则可能漏掉尾部丢失。代理实验需要同时记录首事件时间、事件总数与连接终态。

这段 Python 代理没有使用 nginx,也没有访问生产网关。它展示的是显式缓冲算法对有限输入的影响,不能据此推导某个代理产品的默认配置。部署链路中的压缩、缓存、缓冲和空闲超时仍需各自验证。

响应已经开始以后的失败

mode=fail 在前三个事件之后令上游失败。本次客户端已经得到 HTTP 200 和 chunked 响应头,也完成了三个事件的解析,继续读取时收到 ConnectionResetError;服务端账本的 failure 为 IllegalStateException。

另一轮同类流失败可以呈现不完整 chunked 消息。TCP reset 与读到 EOF 后识别消息截断,取决于服务器关闭方式和客户端读取时序。验收关注响应未完整结束、上游异常被记录及流释放,不将某一个 Python 异常类型当成协议保证。

自动重连如果反复携带同一个最后游标,会反复遇到同一条坏事件。服务端需要决定是保留失败、产生可识别的错误事件、跳过特定记录还是终止订阅;随意跳过会改变业务语义。当前实验直接中止,不伪造一条“成功”的业务事件掩盖异常。

重跑与练习

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

从第00篇取得累计工程并配置 JDK 21,在 play-lab 目录执行:

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

历史隔离证据在 evidence/batch18-21/isolated/:run-final-io/observations.json 的 sse 与 localProxyOnly 保存事件字段、游标状态码和时间;junit/TEST-StreamLabTest.xml 包含事件编码、游标规则与应用停止测试。本轮隔离工程 36 项 JUnit 与 stage 通过,不代表另一次共享工程构建的结果。

反例题:Source 为每条事件生成 id,但发送完马上丢弃对象。客户端重新提交旧 id 后,服务器能否凭这个 id 恢复 payload?不能;标识与保存的内容是两个独立条件。

改动练习:把保留集合改为最多 8 条的滑动日志,在生成第 121 条后删除最旧项。分别以“刚好仍可重放”“比最旧项早一步”“超过最新项”的游标连接,写出明确断言。并发生成与读取时,应从同一受保护的快照计算最旧位置和待发送集合,避免先通过游标校验、随后记录已经被淘汰。

另一个练习是在消费者处理第三条事件之后、保存游标之前模拟崩溃。重连后记录重复事件,并用业务去重键验证效果只发生一次。这个练习验证消费者恢复,不能用服务端给事件编号的单元测试代替。

上一篇:HTTP流式响应。下一篇:WebSocket生命周期。