Hyperf 本身不内置 ETL 引擎,但可通过 flow-php/etl + hyperf/database + hyperf/async-queue 实现千万级日志实时聚合;需规避 chunk 同步阻塞与连接复用问题,改用 Generator 流式读取、协程连接池及 FETCH_NUM/FETCH_ASSOC;rows() 中仅做轻量处理,Loader 层用 insertAsBatch() 批量 upsert;启用 SWOOLE_HOOK_ALL 确保协程化;结合延迟队列实现可重试有序缓冲;Extractor/Loader 时区须显式设为 UTC 防窗口错位。

Hyperf 本身不内置 ETL 引擎,但用 flow-php/etl + hyperf/database + hyperf/async-queue 组合,能跑通千万级日志的实时聚合链路——关键不在框架多强,而在协程调度、内存控制和 Loader 写入节奏是否对齐。
为什么不能直接用 Hyperf/database::chunk 做实时聚合
因为 chunk 是同步阻塞式分页,每次查完一批就等写入完成才查下一批,在高并发写入场景下容易卡住协程调度器;更严重的是,它默认不释放 PDO 连接,协程间复用连接时可能触发 MySQL 的 Packet is bigger than max_allowed_packet 或连接超时。
- 真实错误现象:
PDOException: SQLSTATE[HY000]: General error: 2006 MySQL server has gone away - 正确做法:改用
flow-php/etl的Generator流式读取,配合hyperf/db的协程连接池自动管理 - 必须显式关闭
fetch_mode的对象映射(避免生成大量 DTO 实例),用FETCH_NUM或FETCH_ASSOC直接拿数组
如何让 flow-php/etl 在 Hyperf 里真正“实时”起来
flow-php/etl 默认是批处理模型,所谓“实时”得靠它管道里的 rows() + 协程并发写入来模拟。重点不是拉得多快,而是压得住、吐得稳。
- 在
rows()回调里别做耗时操作(比如 HTTP 请求、文件写入),只做字段映射、类型转换、简单过滤 - 聚合逻辑(如
COUNT/SUM分组)必须下沉到 Loader 层,用INSERT ... ON DUPLICATE KEY UPDATE或REPLACE INTO批量 upsert,避免在 PHP 层维护状态 - Loader 必须用
hyperf/db的insertAsBatch(),且$batchSize控制在 500–2000 行之间——太小吞吐低,太大触发 MySQLmax_allowed_packet - 记得在 Flow::setUp() 前调用
Co::set(['hook_flags' => SWOOLE_HOOK_ALL]),否则某些底层驱动(如 PDO MySQL)不会协程化
延迟队列 + 聚合补偿怎么和 ETL 链路对齐
纯实时链路扛不住抖动,必须加一层“可重试+有序”的缓冲。RabbitMQ 延迟插件或 NSQ 的 REQ 是更稳妥的选择,Kafka 在 Hyperf 里做精确一次语义成本太高。
- ETL 流程中每个
write()成功后,发一条带业务主键的延迟消息(比如 30s 后触发校验),避免重复聚合 - 补偿 Job 里不要重新跑全量 ETL,而是查
SELECT COUNT(*) FROM raw_log WHERE processed = 0 AND created_at > ?,只补漏 - 关键字段如
event_id或trace_id必须全程透传,不能在map()中丢弃,否则补偿无法定位数据源 - 别依赖 RabbitMQ 的死信队列自动重试——Hyperf 的
retry_after和 AMQP 的x-message-ttl容易冲突,统一用应用层重试 + 指数退避
最易被忽略的一点:ETL 管道里所有 Extractor 和 Loader 的时区必须显式设为 UTC,Hyperf 默认用系统时区,而数据中台下游(如 ClickHouse、StarRocks)通常强制 UTC 时间戳,差 8 小时会导致聚合窗口错位。别信文档说“自动适配”,自己在 CSVExtractor 或 DatabaseExtractor 初始化时加 ->withTimezone(new \DateTimeZone('UTC'))。


















