本文目录导读:

在PHP项目中实现实时ETL(Extract, Transform, Load)主要有以下几种方案,具体选择取决于数据源、实时性要求和系统规模:
基于消息队列的流式处理(推荐)
架构设计
// 1. 数据提取 - 使用消息生产
class DataProducer {
private $queue;
public function __construct() {
$this->queue = new RedisQueue(); // 或 RabbitMQ/Kafka
}
public function extract($dataSource) {
// 从数据源提取增量数据
$data = $this->fetchIncrementalData($dataSource);
// 发送到消息队列
foreach ($data as $record) {
$this->queue->publish('etl_pipeline', json_encode($record));
}
}
}
// 2. 数据转换 - 消费者处理
class DataTransformer {
public function transform($message) {
$data = json_decode($message, true);
// 数据清洗和转换
$transformed = [
'id' => (int)$data['id'],
'name' => trim($data['name']),
'amount' => (float)$data['amount'],
'created_at' => date('Y-m-d H:i:s', strtotime($data['timestamp'])),
'status' => $this->normalizeStatus($data['status'])
];
// 业务规则验证
if ($this->validate($transformed)) {
return $transformed;
}
return null;
}
}
// 3. 数据加载 - 批量写入
class DataLoader {
public function load($transformedData) {
// 批量加载到目标数据库
$batchSize = 1000;
$batch = [];
foreach ($transformedData as $record) {
$batch[] = $record;
if (count($batch) >= $batchSize) {
$this->bulkInsert($batch);
$batch = [];
}
}
// 处理剩余数据
if (!empty($batch)) {
$this->bulkInsert($batch);
}
}
}
基于PHP定时任务的微批次处理
// ETL处理脚本 (cron每分钟执行)
class MicroBatchETL {
private $lastProcessedTime;
public function process() {
// 1. 提取增量数据
$this->lastProcessedTime = $this->getLastCheckpoint();
$newData = $this->extractIncremental($this->lastProcessedTime);
if (empty($newData)) {
return;
}
// 2. 内存中转换
$transformed = [];
foreach ($newData as $record) {
$transformed[] = $this->transform($record);
}
// 3. 批量加载
$this->loadBatch($transformed, 'dw_data');
// 4. 更新检查点
$this->updateCheckpoint();
}
private function extractIncremental($lastTime) {
// 使用数据库增量查询
return DB::query("SELECT * FROM source_table
WHERE updated_at > ?", [$lastTime]);
}
private function transform($record) {
// 数据转换逻辑
return [
'id' => $record['id'],
'dimensions' => $this->enrichDimensions($record),
'measures' => $this->calculateMeasures($record)
];
}
private function loadBatch($data, $table) {
// 使用批量插入或UPSERT
DB::table($table)->upsert($data, ['id']);
}
}
集成Change Data Capture (CDC)
// MySQL binlog监听实现
class BinlogETLConsumer {
public function handleEvent($event) {
switch ($event->getAction()) {
case 'insert':
$this->handleInsert($event->getRow());
break;
case 'update':
$this->handleUpdate($event->getRow());
break;
case 'delete':
$this->handleDelete($event->getRow());
break;
}
}
private function handleInsert($row) {
// 实时转换
$transformed = $this->transformToDW($row);
// 立即加载到数据仓库
$this->loader->insert($transformed);
// 同时更新缓存
$this->cache->set("data:{$transformed['id']}", $transformed);
}
}
高性能PHP-FPM结合ReactPHP
// 异步非阻塞处理
use React\EventLoop\Factory;
use React\Stream\ReadableResourceStream;
class AsyncRealtimeETL {
public function run() {
$loop = Factory::create();
// 建立持久连接
$sourceStream = new ReadableResourceStream($this->getSourceConnection(), $loop);
$outputStream = new WritableResourceStream($this->getTargetConnection(), $loop);
// 管道处理
$sourceStream->pipe(new TransformStream())->pipe($outputStream);
$loop->run();
}
}
// 自定义转换流
class TransformStream extends ThroughStream {
public function filter($data) {
$decoded = json_decode($data, true);
$transformed = $this->transform($decoded);
return json_encode($transformed) . "\n";
}
private function transform($data) {
// 实时转换逻辑
return [
'id' => $data['id'],
'value' => $data['value'] * 1.2,
'timestamp' => date('Y-m-d H:i:s')
];
}
}
完整实现示例
<?php
class RealTimeETLManager {
private $extractor;
private $transformer;
private $loader;
private $logger;
public function __construct() {
$this->extractor = new KafkaExtractor('source_topic');
$this->transformer = new DataTransformer();
$this->loader = new ClickHouseLoader();
$this->logger = new MonologLogger('etl');
}
public function start() {
// 启动消费者组
$this->extractor->consume(function($message) {
try {
// 开始计时
$startTime = microtime(true);
// 数据转换
$transformed = $this->transformer->transform($message);
if ($transformed !== null) {
// 数据加载
$this->loader->loadIncremental($transformed);
// 记录延迟
$latency = (microtime(true) - $startTime) * 1000;
$this->logger->info('ETL latency', [
'id' => $transformed['id'],
'latency_ms' => $latency
]);
}
} catch (Exception $e) {
$this->logger->error('ETL process failed', [
'error' => $e->getMessage(),
'data' => $message
]);
// 发送到死信队列
$this->extractor->sendToDLQ($message);
}
});
}
}
性能优化建议
批次合并
class BatchBuffer {
private $buffer = [];
private $maxSize = 5000;
private $flushInterval = 2; // 秒
public function add($record) {
$this->buffer[] = $record;
if (count($this->buffer) >= $this->maxSize) {
$this->flush();
}
}
public function flush() {
if (empty($this->buffer)) return;
// 批量写入
$this->loader->bulkInsert($this->buffer);
$this->buffer = [];
}
}
使用连接池
class ConnectionPool {
private $pool = [];
private $maxConnections = 10;
public function getConnection() {
// 从池中获取或创建新连接
return $this->pool[] ?? $this->createConnection();
}
}
选择建议
| 场景 | 推荐方案 | 优势 |
|---|---|---|
| 高吞吐、低延迟 | Kafka + PHP消费者组 | 可扩展、持久化 |
| 简单实时同步 | MySQL binlog + CDC | 对源系统无侵入 |
| 资源受限 | 微批次处理 (cron) | 实现简单 |
| 复杂转换 | Redis队列 + 工作进程 | 灵活、易调试 |
监控与运维
// 健康检查端点
class ETLHealthCheck {
public function check() {
return [
'queue_lag' => $this->getQueueLag(),
'throughput' => $this->getThroughput(),
'error_rate' => $this->getErrorRate(),
'last_success_time' => $this->getLastSuccessTime()
];
}
}
对于生产环境,建议:
- 使用 Apache Kafka 或 Redis Streams 作为消息中间件
- 结合 Supervisor 管理PHP进程
- 实现断点续传和死信队列
- 监控ETL延迟和数据一致性