goletter / hyperf-sse
Server-Sent Events (SSE) component for Hyperf with Redis Stream, connection tickets, and optional presence.
Requires
- php: >=8.1
- hyperf/codec: ~3.1.0
- hyperf/context: ~3.1.0
- hyperf/contract: ~3.1.0
- hyperf/coroutine: ~3.1.0
- hyperf/engine: ~2.0 || ~3.0
- hyperf/http-server: ~3.1.0
- hyperf/redis: ~3.1.0
- hyperf/support: ~3.1.0
- psr/container: ^1.0 || ^2.0
- psr/http-message: ^1.0 || ^2.0
Requires (Dev)
- friendsofphp/php-cs-fixer: ^3.0
- mockery/mockery: ^1.0
- phpstan/phpstan: ^1.0
- phpunit/phpunit: ^10.0
- swoole/ide-helper: dev-master
README
面向 Hyperf 3.1 的 Server-Sent Events(SSE)扩展包。
包名:
goletter/hyperf-sse· PHP>= 8.1· Hyperf~3.1.0
SSE 是单向通信:只从服务端推到浏览器。前端提交任务仍走普通 HTTP API;任务进度 / 完成通知用 SSE 推回来。
一、整体怎么用(先看这个)
你们的登录是 Authorization: Bearer ...,而浏览器原生 EventSource 不能自定义 Header,所以标准用法是三步:
① 前端带 Authorization 调 POST /sse/ticket
↓ 拿到短时一次性 ticket
② 前端 new EventSource('/sse/connect?ticket=xxx')
↓ 服务端把 ticket 换成 userId,开始挂起推流
③ 业务里 Sse::to($userId)->event(...)->data(...)->send()
↓ 只有「已连接且 subject 相同」的浏览器收到事件
对应关系(务必一致):
| 步骤 | subject(通常是用户 ID) |
|---|---|
| 发 ticket | issueTicket(42) → ticket 绑定用户 42 |
| 建立连接 | ticket 消费后变成 subscribe(42) |
| 业务推送 | Sse::to(42)->...->send() |
包不会自动知道该推谁。推给谁,由业务数据里的
user_id决定(例如任务所有者)。
二、安装
composer require goletter/hyperf-sse php bin/hyperf.php vendor:publish goletter/hyperf-sse
本地联调(path):
{
"repositories": [
{
"type": "path",
"url": "../hyperf-sse",
"options": { "symlink": true }
}
],
"require": {
"goletter/hyperf-sse": "dev-main"
}
}
发布后得到 config/autoload/sse.php。
.env 最少配置:
SSE_DRIVER=redis SSE_REDIS_POOL=default SSE_HEARTBEAT_INTERVAL=15 SSE_MAX_IDLE_TIME=7200 SSE_TICKET_ENABLED=true SSE_TICKET_TTL=60
鉴权中间件里请把登录用户 ID 写到请求属性(默认名 authenticated_user_id,可在配置里改 subject_attribute)。
三、后端:两个接口
完整可复制示例:examples/SseController.php。
1)发 ticket(需要登录 / Authorization)
use Goletter\Sse\SseService; use Hyperf\HttpServer\Annotation\Middleware; use Hyperf\HttpServer\Annotation\PostMapping; #[PostMapping(path: '/sse/ticket')] #[Middleware(YourAuthMiddleware::class)] // 校验 Authorization,写入 authenticated_user_id public function ticket(SseService $sse): array { $userId = $this->request->getAttribute('authenticated_user_id'); return [ 'ticket' => $sse->issueTicket($userId), // 默认 60 秒,一次性 'ttl' => 60, ]; }
2)建立 SSE 长连接(用 ticket,一般不再校验 Authorization)
use Hyperf\HttpServer\Annotation\GetMapping; #[GetMapping(path: '/sse/connect')] public function connect(SseService $sse): void { // 内部会:读 ?ticket= → 换成 userId → 挂起推流 // 返回类型必须是 void,不要再 return json $sse->connect($this->request); }
connect() 解析 subject 的顺序:
- 查询参数
?ticket=(推荐) - 否则读请求属性
authenticated_user_id(Cookie/Session 场景)
四、前端:怎么连、怎么听
前端不用做心跳。服务端会定时发 : ping 保活链路;这种注释帧浏览器不会抛给 JS。EventSource 断线后会自动重连。
async function openSse() { const token = localStorage.getItem('token'); // 你们现有的登录 Token // ① 用 Authorization 换 ticket const res = await fetch('/sse/ticket', { method: 'POST', headers: { Authorization: `Bearer ${token}`, Accept: 'application/json', }, }); if (!res.ok) throw new Error('获取 SSE ticket 失败'); const { ticket } = await res.json(); // ② 建立 SSE(这里不能带 Authorization,所以用 ticket) const es = new EventSource(`/sse/connect?ticket=${encodeURIComponent(ticket)}`); // ③ 听业务事件(事件名要和后端 event() 一致) es.addEventListener('task_completed', (e) => { const data = JSON.parse(e.data); console.log('任务完成', data); // 建议再调一次任务详情 API,用数据库状态校准 }); // 服务端主动关闭(空闲超时等) es.addEventListener('server_close', () => { es.close(); }); // 网络异常:EventSource 会自己重连;这里只做提示即可 es.onerror = () => { console.warn('SSE 连接异常,浏览器将自动重连'); }; return es; } // 页面进入时打开,离开时关闭 const es = await openSse(); // window.addEventListener('beforeunload', () => es.close());
五、业务侧:怎么推给指定用户
任务完成、导出完成等场景——你已经知道目标用户 ID:
use Goletter\Sse\Sse; // 推荐:链式 Sse::to($task->user_id) // 推给谁 ->event('task_completed') // 前端 addEventListener 的名字 ->data([ // 前端 JSON.parse(e.data) 'task_id' => $task->id, 'status' => 'completed', ]) ->send(); // 推给多人 Sse::toMany([1, 2, 3])->event('notice')->data(['msg' => 'hello'])->send(); // 所有在线连接都收(广播) Sse::broadcast('maintenance', ['message' => '系统将在 10 分钟后维护']); // 助手函数(等价写法) sse_notify($task->user_id, 'task_completed', ['task_id' => $task->id]); sse_broadcast('maintenance', ['message' => '...']);
注入写法:
public function __construct(private SseService $sse) {} $this->sse->to($userId)->event('notice')->data('hi')->send();
完整时序(任务通知)
用户 42 打开页面
→ POST /sse/ticket(Authorization)→ ticket
→ EventSource(/sse/connect?ticket=) → 服务端 subscribe(42)
用户 42 提交导出任务
→ POST /api/export(普通 API)
→ Worker 跑完,更新 DB status=completed
→ Sse::to(42)->event('task_completed')->data(...)->send()
用户 42 浏览器
→ 收到 task_completed
→ 再 GET /api/tasks/1001 校准状态
用户当时没连着 SSE:notify 仍可写入 Redis Stream(有限保留);重连后可能靠 Last-Event-ID 补到,仍应以查库为准。
六、心跳与断线(前后端分别做什么)
| 端 | 要不要心跳 | 说明 |
|---|---|---|
| 服务端 | 要(已内置) | 默认每 15s 发 : ping,防代理掐空闲连接 |
| 前端 | 不用 | 收不到 comment;不必发 ping,也不必检测心跳 |
| 浏览器 | 自动重连 | EventSource 断线后会重连并带 Last-Event-ID |
空闲过久(默认 7200s 无业务消息)服务端会发 server_close 并结束连接,前端监听到后 es.close() 即可。
七、Nginx(必配,否则容易卡住)
location /sse/ { proxy_pass http://hyperf_upstream; proxy_http_version 1.1; proxy_set_header Connection ''; proxy_buffering off; proxy_cache off; gzip off; proxy_read_timeout 2h; proxy_send_timeout 2h; }
心跳间隔必须 小于 链路里最短的空闲超时。
八、可选:在线 Presence
默认关闭。需要「谁在线」时再开:
SSE_PRESENCE_ENABLED=true
Sse::online()->isOnline(42); Sse::online()->list();
推送不依赖在线列表:即使不开 Presence,Sse::to(42) 照样能推。
九、API 速查
| 方法 | 用途 |
|---|---|
$sse->issueTicket($userId) |
登录态换短时连接票据 |
$sse->connect($request) |
建立 SSE 推流(Controller 返回 void) |
Sse::to($id)->event()->data()->send() |
推给一个用户 |
Sse::toMany([$id, ...])->...->send() |
推给多个用户 |
Sse::broadcast($event, $data) |
广播 |
Sse::notify($id, $event, $data) |
底层定向推送 |
Sse::online() |
Presence(需开启配置) |
sse() / sse_to() / sse_notify() |
助手函数 |
十、注意点
- 不要让前端传
user_id来订阅;用 ticket 或服务端登录态。 - 不要把长期 JWT 放在 URL 里;只用短时一次性 ticket。
- SSE 只推小消息(任务 ID、状态);大结果走普通下载 / 查询接口。
connect的 Controller 必须void,调用后不要再返回 JSON。- 多端同时打开:同一
userId多个连接都会收到同一次to($userId)。