服务器发送事件(SSE)

OxPHP 通过服务器发送事件(Server-Sent Events)协议向客户端流式推送实时数据,并内置背压机制。只需在 PHP 脚本中设置 Content-Type: text/event-stream 并调用 oxphp_stream_flush(),其余的交给 OxPHP 处理。

工作原理

一个流式响应会经历一个可预测的生命周期:

  1. 设置响应头。 你的 PHP 脚本通过 header() 设置 Content-Type: text/event-stream,并使用 echo 写入符合 SSE 格式的文本行。
  2. 首次 flush 打开流。 第一次调用 oxphp_stream_flush() 会把 HTTP 响应头发送给客户端并进入流式模式。客户端连接保持打开。
  3. 每次 flush 发送一个数据块。 之后每次调用 oxphp_stream_flush() 都会把缓冲的输出作为一个新的数据块刷出,立即投递给客户端。
  4. 缓冲区填满时触发背压。 OxPHP 在 PHP 工作进程与客户端之间维护一个最多 64 个数据块的内部缓冲区。当缓冲区已满时——因为慢速客户端尚未消费掉先前的数据块——oxphp_stream_flush() 会阻塞,直到腾出空间为止。这可以防止内存无限增长。
  5. 退出或断开时清理。 当 PHP 脚本执行完毕时,OxPHP 会优雅地关闭连接。如果客户端在流中途断开,OxPHP 会在下一次 flush 时检测到已关闭的通道,将 PHP 的 connection_aborted() 标志置为 true,并准备一次优雅退出——检查 connection_aborted() 的可移植循环会沿其正常终止路径干净地退出,而不检查它的循环仍会在下一次 flush 时被一次隐式退出所终止。
Note

保持事件负载(请求体)较小,以维持平滑的吞吐量。较大的负载会很快填满这个 64 块的缓冲区,导致 PHP 在每次 flush 时阻塞。

示例

基础 SSE 流

php
<?php header('Content-Type: text/event-stream'); header('Cache-Control: no-cache'); header('Connection: keep-alive'); for ($i = 0; $i < 100; $i++) { $data = json_encode(['counter' => $i, 'time' => microtime(true)]); echo "id: {$i}\n"; echo "event: tick\n"; echo "data: {$data}\n\n"; oxphp_stream_flush(); sleep(1); // Send a comment heartbeat every 15 seconds to keep proxies from closing idle connections if ($i % 15 === 0) { echo ": heartbeat\n\n"; oxphp_stream_flush(); } }

检查流式状态

使用 oxphp_is_streaming() 来检查当前请求是否已经处于流式模式。这在中间件或共享的请求处理器中很有用:

php
<?php if (!oxphp_is_streaming()) { header('Content-Type: text/event-stream'); header('Cache-Control: no-cache'); } echo "data: {\"status\": \"connected\"}\n\n"; oxphp_stream_flush();

检测客户端断开

长时间运行的 SSE 循环应当检查 connection_aborted(),以便在客户端关闭连接时干净地跳出。这与标准的 PHP / php-fpm 惯用法一致,并让脚本在退出前运行任何清理逻辑(关闭数据库句柄、释放锁、完成 finally 块):

php
<?php header('Content-Type: text/event-stream'); header('Cache-Control: no-cache'); $db = new PDO(/* ... */); try { while (!connection_aborted()) { echo "data: " . json_encode(['ts' => time()]) . "\n\n"; oxphp_stream_flush(); sleep(1); } } finally { $db = null; // runs on normal exit AND on connection_aborted exit }
Note

如果脚本从不检查 connection_aborted(),OxPHP 仍会在客户端断开后的下一次 flush 时通过一次隐式退出终止它,但那些绕过 flush 调用的代码路径之后的 finally 块可能不会运行。对于持有外部资源的代码,请优先使用显式检查。

使用原生 flush()

PHP 的原生 flush() 也可用于流式传输,但需要先清空所有输出缓冲层。请优先使用 oxphp_stream_flush()——它会自动管理输出缓冲,并与 OxPHP 的背压系统集成。

php
<?php header('Content-Type: text/event-stream'); header('Cache-Control: no-cache'); while (ob_get_level()) { ob_end_clean(); } for ($i = 0; $i < 100; $i++) { echo "data: " . json_encode(['counter' => $i]) . "\n\n"; flush(); sleep(1); }

工作进程模式下的 SSE

SSE 在标准模式和工作进程模式下都能工作。在工作进程模式下,流式连接会在整个流的持续时间内占用该工作进程。只有在脚本执行完毕后,工作进程才会处理下一个请求。

php
<?php require __DIR__ . '/../vendor/autoload.php'; $redis = new Redis(); $redis->pconnect('redis', 6379); oxphp_worker(function () use ($redis) { if (($_SERVER['HTTP_ACCEPT'] ?? '') !== 'text/event-stream') { http_response_code(400); echo json_encode(['error' => 'SSE only']); return; } header('Content-Type: text/event-stream'); header('Cache-Control: no-cache'); while (true) { $message = $redis->brPop('events', 25); if ($message) { echo "data: {$message[1]}\n\n"; } else { // No message within timeout — send heartbeat to keep the connection alive echo ": heartbeat\n\n"; } oxphp_stream_flush(); } });

故障排查

脚本结束前客户端收不到任何数据

PHP 的输出缓冲正在捕获输出,而不是将其流式发出。这发生在 OB 层处于激活状态且未调用 oxphp_stream_flush() 时。

修复方法: 在每个事件之后调用 oxphp_stream_flush()。该函数会刷出所有 PHP 输出缓冲层,并把累积的输出作为一个数据块发送。

SSE 连接在几分钟后被关闭

PHP 的 max_execution_time 触发并终止了脚本。SSE 流的运行时间必须超过所配置的限制。

修复方法: 在流式脚本的顶部禁用按请求计算的执行计时器:

php
set_time_limit(0);

这优于全局设置 max_execution_time = 0——它会为非 SSE 端点保留该限制。或者,如果整个实例都专门用于长时间运行的流:

php.ini
; php.ini max_execution_time = 0
中间代理关闭空闲的 SSE 连接

负载均衡器和代理通常会关闭 30–60 秒内没有数据传输的连接。

修复方法: 定期发送一个注释心跳,以保持连接活跃:

php
echo ": heartbeat\n\n"; oxphp_stream_flush();
oxphp_stream_flush() 返回 false

在同一个请求中较早地调用了 oxphp_finish_request()。一旦响应已经结束,就无法再进行流式传输。请检查你的代码,看是否在流式传输开始之前无意中调用了 oxphp_finish_request()

Docker 示例

SSE 端点要求禁用 PHP 的执行计时器,或将其设置得较高。每个活跃的 SSE 连接都会在整个流的持续时间内占用一个 PHP 工作进程,因此要根据你预期的并发流数量来调整工作进程池的大小。

compose.yaml
services: app: image: ghcr.io/oxphp/oxphp:0.10.0 ports: - "8080:8080" volumes: - ./src:/var/www/html environment: DOCUMENT_ROOT: "/var/www/html/public" ENTRY_FILE: "index.php" PHP_WORKERS: "32"

每个流式脚本都应在顶部调用 set_time_limit(0),这样按请求计算的计时器就不会在流中途触发。这能让全局的 max_execution_time 对非 SSE 请求继续生效。

最佳实践

  • 使用 oxphp_stream_flush() 而非原生 flush(),以获得自动的输出缓冲管理和背压集成。
  • 每 20–30 秒发送一次周期性的注释心跳: heartbeat\n\n),以防止中间代理关闭空闲连接,并尽早检测客户端断开。
  • 保持事件负载(请求体)较小。 较大的负载会更快填满这个 64 块的缓冲区,导致 PHP 在每次 flush 时停顿。对于大数据,应发送一个事件 ID,让客户端通过单独的请求去获取完整负载。
  • 为长时间运行的 SSE 端点按脚本禁用 PHP 的执行计时器,使用 set_time_limit(0),或者把 max_execution_time 设置得足够高,以覆盖你预期的最长流持续时间。
  • 根据并发流的峰值来调整工作进程池的大小。 每个活跃的 SSE 连接在其整个持续时间内都会占用一个 PHP 工作进程。至少要为每个预期的并发客户端预留一个工作进程,另外再为常规的非 SSE 请求预留额外的工作进程。

注意事项

  • 对于流式响应,Brotli 压缩会被自动跳过。压缩只适用于完全缓冲的响应。
  • 如果在同一个请求上已经调用过 oxphp_finish_request(),则 oxphp_stream_flush() 返回 false
  • 在工作进程模式下,工作进程会在整个流的持续时间内保持被占用,只有在 PHP 脚本退出后才处理下一个请求。

关闭时的行为

当服务器收到 SIGTERM 时(滚动部署、docker stop、Kubernetes pod 驱逐),它会优雅地进行排空:

  • 每个打开的 SSE 流都会在其下一次 flush 时被干净地结束——其处理器会像连接已关闭那样退出,因此 register_shutdown_function() 回调仍会运行,且 error_get_last()['message'] 读取到 Request cancelled (shutdown)
  • HTTP/2 客户端会收到一个 GOAWAY 帧;HTTP/1.1 keep-alive 连接会被关闭。浏览器的 EventSource 会自动重连——当负载均衡器位于集群前面时,会重连到一个健康的实例。
  • 该流不会收到 503:它的 200 响应头已经发送,因此状态码无法被改写。应将客户端设计为从最后一个事件 id 处恢复(随每个事件发送 id:;重连时浏览器会将其作为 Last-Event-ID 请求头回传,在 PHP 中可通过 $_SERVER['HTTP_LAST_EVENT_ID'] 获取)。

在 SIGTERM 时正在处理中的普通(非流式)请求不会被中断:它们会获得整个排空窗口来正常完成,只有那些在 DRAIN_TIMEOUT_SECONDS(默认 25)到期时仍在运行的请求才会被取消,另有约 2 秒的时间来收尾。请将编排器的终止宽限期设置为高于 DRAIN_TIMEOUT_SECONDS + 2。

这里所说的“流式”指的是任何已经刷出过分块输出的响应——服务器无法区分一个有限的流式下载和一个无限的事件流,因此在 SIGTERM 时正在处理中的大型已刷出下载也会被提前结束,而不仅仅是 SSE。一个在 SIGTERM 之前就调用过 oxphp_finish_request() 的请求会被视为普通请求:它的响应已经完成,其剩余的后台工作会获得排空窗口。

一个阻塞在单个长时间运行的原生调用内部的处理器(原生 sleep()、繁重的 preg_match、阻塞式数据库查询)在该调用返回之前无法被中断。在工作进程模式下,流式循环中应优先使用协作式的 oxphp_sleep()——它会让出给纤程调度器,并在关闭时立即被唤醒。在工作进程模式之外,oxphp_sleep() 会退化为常规的阻塞式 sleep,因此流会在 sleep 返回后的下一次 flush 时对关闭做出反应。

参见

  • 工作进程模式 —— 持久化的 PHP 进程,用于降低启动开销
  • 超时 —— 为长时间运行的连接配置或禁用请求超时
  • PHP 函数 —— oxphp_stream_flush()oxphp_is_streaming() 的完整参考
  • 压缩 —— Brotli 压缩行为以及哪些响应会被压缩