本文目录导读:

在PHP项目中实现ETL(Extract, Transform, Load),通常有以下几种主流方案,具体选择取决于数据规模、复杂度以及性能要求。
纯PHP脚本实现(适合小规模、简单场景)
通过编写PHP脚本,手动完成数据抽取、转换和加载三个步骤。
示例代码
<?php
// 1. Extract - 从CSV文件抽取数据
$sourceFile = 'data.csv';
$data = [];
if (($handle = fopen($sourceFile, 'r')) !== false) {
while (($row = fgetcsv($handle)) !== false) {
$data[] = $row;
}
fclose($handle);
}
// 2. Transform - 数据转换
$transformed = [];
foreach ($data as $index => $row) {
if ($index == 0) continue; // 跳过表头
$transformed[] = [
'name' => trim($row[0]),
'age' => (int)$row[1],
'email' => filter_var($row[2], FILTER_VALIDATE_EMAIL)
];
}
// 3. Load - 加载到MySQL
$pdo = new PDO('mysql:host=localhost;dbname=test', 'user', 'pass');
$stmt = $pdo->prepare('INSERT INTO users (name, age, email) VALUES (?, ?, ?)');
foreach ($transformed as $record) {
$stmt->execute([$record['name'], $record['age'], $record['email']]);
}
?>
缺点:大数据量时内存消耗大,性能差,缺乏容错机制。
使用ETL专用PHP库
如果不想完整开发,可以借助开源PHP ETL库。
推荐库
- Flow\ETL(PHP 8.1+)
- Keboola PHP ETL
- Earmon/ETL(轻量)
示例(Flow\ETL)
use Flow\ETL\ETL;
use Flow\ETL\Extractor\CSVExtractor;
use Flow\ETL\Transformer\CallbackTransformer;
use Flow\ETL\Loader\PDOLoader;
ETL::extract(new CSVExtractor('data.csv'))
->transform(new CallbackTransformer(fn($row) => [
'name' => trim($row['name']),
'age' => (int)$row['age'],
'email' => strtolower($row['email'])
]))
->load(new PDOLoader($pdo, 'users'))
->run();
优点:代码简洁,支持管道式处理,可分批处理大文件。
结合队列与多进程(处理大数据量)
当数据量大到单进程处理太慢时,可以用队列(RabbitMQ/Redis)+ 多进程/Workerman 并行处理。
流程
- Extract:读取数据源,分割成小批次,推入队列。
- Transform:多个worker从队列消费,并行转换。
- Load:worker将结果写入目标数据库/文件。
示例(使用Redis队列 + 多进程)
// Producer - 读取并推送
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
$handle = fopen('bigdata.csv', 'r');
while (($row = fgetcsv($handle)) !== false) {
$redis->lPush('etl_queue', json_encode($row));
}
fclose($handle);
// Worker(多进程运行)
$redis = new Redis();
$redis->connect('127.0.0.1', 6379);
$pdo = new PDO('mysql:host=localhost;dbname=test', 'user', 'pass');
while ($data = $redis->brPop('etl_queue', 5)) {
$row = json_decode($data[1], true);
// 转换
$transformed = transformRow($row);
// 加载
$stmt = $pdo->prepare('INSERT INTO users (...) VALUES (...)');
$stmt->execute($transformed);
}
适用场景:每天百万级+数据,需要高吞吐量。
调用外部ETL工具(推荐生产环境)
PHP不需要自己做ETL,可以通过系统调用或API调用专业ETL工具,PHP只做任务调度和结果处理。
常用工具
- Apache NiFi(REST API)
- Talend Open Studio(命令行)
- Pentaho Data Integration(命令行)
- Airflow(Python,可调度PHP脚本)
示例(PHP调用NiFi API)
$ch = curl_init('http://nifi-server:8080/nifi-api/processors/.../run-status');
curl_setopt($ch, CURLOPT_POST, true);
curl_setopt($ch, CURLOPT_POSTFIELDS, json_encode(['state' => 'RUNNING']));
curl_setopt($ch, CURLOPT_HTTPHEADER, ['Content-Type: application/json']);
curl_exec($ch);
优点:ETL逻辑由专业工具处理,PHP只做编排,稳定可靠。
选择建议
| 场景 | 推荐方案 |
|---|---|
| 数据量小(<10万条),简单转换 | 纯PHP或Flow\ETL |
| 数据量较大,需要管道处理 | Flow\ETL / 队列+多进程 |
| 数据量巨大(GB/TB级),复杂ETL | 调用外部ETL工具(NiFi/Talend) |
| 需要在Web界面监控ETL任务 | Flow\ETL + Symfony/Laravel命令 |
| 实时/准实时ETL(CDC) | PHP + Kafka + 消费者进程 |
生产环境最佳实践:
- 不要用PHP做复杂SQL级的Join/聚合(PHP内存和CPU不适合);应尽量在数据库层面(SQL/存储过程)完成,PHP只做搬运。
- 大文件必须用流式读取(
fgetcsv配合yield),避免全量加载到内存。 - 增加断点续传机制(记录已处理的行号或ID)。
- 使用事务批量提交(每1000行commit一次),提升写入性能。
如果你能提供更具体的数据源类型和目标数据库,我可以给出更精确的代码示例。