本文目录导读:

在 PHP 中编写 ETL(Extract-Transform-Load)过程,有多种方法,从简单的脚本到使用专门的库,这里提供几种常见且实用的实现方式。
基础手动实现(适合小规模数据)
<?php
class SimpleETL {
private $source;
private $destination;
private $transformations = [];
public function __construct($source, $destination) {
$this->source = $source;
$this->destination = $destination;
}
public function addTransformation(callable $transform) {
$this->transformations[] = $transform;
return $this;
}
public function extract() {
// 从源提取数据
$data = [];
$handle = fopen($this->source, 'r');
while (($row = fgetcsv($handle)) !== false) {
$data[] = $row;
}
fclose($handle);
return $data;
}
public function transform($data) {
// 应用所有转换
foreach ($this->transformations as $transform) {
$data = $transform($data);
}
return $data;
}
public function load($data) {
// 加载到目标
$handle = fopen($this->destination, 'w');
foreach ($data as $row) {
fputcsv($handle, $row);
}
fclose($handle);
}
public function run() {
echo "开始ETL过程...\n";
echo "1. 提取数据...\n";
$data = $this->extract();
echo "2. 转换数据...\n";
$transformed = $this->transform($data);
echo "3. 加载数据...\n";
$this->load($transformed);
echo "ETL完成!\n";
}
}
// 使用示例
$etl = new SimpleETL('input.csv', 'output.csv');
$etl->addTransformation(function($rows) {
return array_filter($rows, function($row) {
return !empty($row[0]); // 过滤空行
});
});
$etl->addTransformation(function($rows) {
return array_map(function($row) {
// 数据转换逻辑
$row[2] = strtoupper($row[2]); // 转大写
$row[4] = (int)$row[4] * 2; // 数值运算
return $row;
}, $rows);
});
$etl->run();
使用 PDO 连接数据库(更实际的场景)
<?php
class DatabaseETL {
private $sourceDB;
private $destDB;
private $batchSize = 1000;
public function __construct($sourceConfig, $destConfig) {
$this->sourceDB = new PDO(
"mysql:host={$sourceConfig['host']};dbname={$sourceConfig['dbname']}",
$sourceConfig['user'],
$sourceConfig['pass']
);
$this->destDB = new PDO(
"mysql:host={$destConfig['host']};dbname={$destConfig['dbname']}",
$destConfig['user'],
$destConfig['pass']
);
$this->sourceDB->setAttribute(PDO::ATTR_ERRMODE, PDO::ERRMODE_EXCEPTION);
$this->destDB->setAttribute(PDO::ATTR_ERRMODE, PDO::ERRMODE_EXCEPTION);
}
public function extract($query) {
echo "Extracting data...\n";
$stmt = $this->sourceDB->query($query);
$data = [];
while ($row = $stmt->fetch(PDO::FETCH_ASSOC)) {
$data[] = $row;
}
return $data;
}
public function transform($data) {
echo "Transforming data...\n";
return array_map(function($row) {
// 转换逻辑示例
if (isset($row['created_at'])) {
$row['created_at'] = date('Y-m-d H:i:s', strtotime($row['created_at']));
}
if (isset($row['email'])) {
$row['email'] = strtolower(trim($row['email']));
}
if (isset($row['status'])) {
$statusMap = [
'active' => 1,
'inactive' => 0,
'pending' => 2
];
$row['status'] = $statusMap[$row['status']] ?? 0;
}
return $row;
}, $data);
}
public function load($tableName, $data) {
echo "Loading data to {$tableName}...\n";
$this->destDB->beginTransaction();
try {
$batchCount = 0;
foreach ($data as $row) {
// 构建插入语句
$columns = implode(', ', array_keys($row));
$placeholders = ':' . implode(', :', array_keys($row));
$sql = "INSERT INTO {$tableName} ({$columns}) VALUES ({$placeholders})";
$stmt = $this->destDB->prepare($sql);
$stmt->execute($row);
$batchCount++;
// 分批次提交
if ($batchCount >= $this->batchSize) {
$this->destDB->commit();
$this->destDB->beginTransaction();
$batchCount = 0;
}
}
$this->destDB->commit();
echo "Loaded " . count($data) . " records successfully.\n";
} catch (Exception $e) {
$this->destDB->rollBack();
throw $e;
}
}
public function run($extractQuery, $tableName) {
try {
$data = $this->extract($extractQuery);
$transformed = $this->transform($data);
$this->load($tableName, $transformed);
echo "ETL process completed successfully!\n";
} catch (Exception $e) {
echo "ETL failed: " . $e->getMessage() . "\n";
throw $e;
}
}
}
// 使用示例
$sourceConfig = [
'host' => 'localhost',
'dbname' => 'source_db',
'user' => 'user',
'pass' => 'password'
];
$destConfig = [
'host' => 'localhost',
'dbname' => 'dest_db',
'user' => 'user',
'pass' => 'password'
];
$etl = new DatabaseETL($sourceConfig, $destConfig);
$query = "SELECT id, name, email, created_at, status FROM users WHERE created_at >= '2023-01-01'";
$etl->run($query, 'dim_users');
使用 ETL 框架(推荐生产环境)
<?php
// 使用 EasyRdf 或 Flow 等框架的示例
use PhpETL\Pipeline\Pipeline;
use PhpETL\Transform\Filter;
use PhpETL\Transform\Map;
use PhpETL\Extract\CSVExtractor;
use PhpETL\Load\DatabaseLoader;
class ProductionETL {
public function execute() {
$pipeline = new Pipeline();
// 配置提取器
$extractor = new CSVExtractor([
'file' => '/path/to/input.csv',
'delimiter' => ',',
'enclosure' => '"',
'escape' => '\\'
]);
// 配置转换器
$transformers = [
new Filter(function($row) {
return isset($row['email']) && filter_var($row['email'], FILTER_VALIDATE_EMAIL);
}),
new Map(function($row) {
return [
'id' => (int)$row['id'],
'email' => strtolower($row['email']),
'full_name' => trim($row['first_name'] . ' ' . $row['last_name']),
'created_at' => date('Y-m-d H:i:s'),
'source' => 'csv_import',
'processed_at' => date('Y-m-d H:i:s')
];
})
];
// 配置加载器
$loader = new DatabaseLoader([
'dsn' => 'mysql:host=localhost;dbname=warehouse',
'username' => 'user',
'password' => 'pass',
'table' => 'users_staging',
'batch_size' => 1000
]);
// 执行管道
$pipeline
->setExtractor($extractor)
->setTransformers($transformers)
->setLoader($loader)
->execute();
}
}
// 运行
$etl = new ProductionETL();
$etl->execute();
高级特性实现
1 错误处理和日志
class ErrorHandlingETL {
private $errors = [];
private $logger;
public function __construct() {
$this->logger = new Logger('etl');
}
public function extractWithRetry($source, $maxRetries = 3) {
for ($attempt = 1; $attempt <= $maxRetries; $attempt++) {
try {
return $this->extract($source);
} catch (Exception $e) {
$this->logger->warning("Extraction attempt {$attempt} failed: " . $e->getMessage());
if ($attempt === $maxRetries) {
throw $e;
}
sleep(5 * $attempt); // 指数退避
}
}
}
public function validateData($data) {
$valid = [];
foreach ($data as $row) {
try {
$valid[] = $this->validateRow($row);
} catch (ValidationException $e) {
$this->errors[] = [
'row' => $row,
'error' => $e->getMessage(),
'timestamp' => date('Y-m-d H:i:s')
];
$this->logger->error("Row validation failed: " . $e->getMessage());
}
}
return $valid;
}
}
2 增量加载
class IncrementalETL {
private $lastRunTimestamp;
private $watermarkTable = 'etl_watermarks';
public function getLastRunTimestamp($jobName) {
$stmt = $this->sourceDB->prepare(
"SELECT last_run FROM {$this->watermarkTable} WHERE job_name = ?"
);
$stmt->execute([$jobName]);
return $stmt->fetchColumn();
}
public function extractIncremental($jobName, $mainQuery) {
$lastRun = $this->getLastRunTimestamp($jobName);
$query = $mainQuery . " AND created_at > ?";
$stmt = $this->sourceDB->prepare($query);
$stmt->execute([$lastRun]);
return $stmt->fetchAll(PDO::FETCH_ASSOC);
}
}
最佳实践建议
- 批处理:始终使用批次处理大数据集
- 内存管理:注意内存使用,对大文件使用流式处理
- 错误处理:实现完整的错误处理和恢复机制
- 日志记录:记录详细的执行日志
- 性能优化:使用索引、批量插入、并行处理
- 可重试性:确保ETL可以安全重试
- 测试:为每个ETL编写单元和集成测试
- 监控:监控执行时间、数据量、错误率
选择哪种方法取决于你的具体需求:
- 小规模简单任务:手动实现
- 中等规模:使用 PDO 自己封装
- 生产级系统:使用成熟的 ETL 框架或工具