PHP ETL过程怎么写

wen PHP项目 3

本文目录导读:

PHP ETL过程怎么写

  1. 基础手动实现(适合小规模数据)
  2. 使用 PDO 连接数据库(更实际的场景)
  3. 使用 ETL 框架(推荐生产环境)
  4. 高级特性实现
  5. 最佳实践建议

在 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);
    }
}

最佳实践建议

  1. 批处理:始终使用批次处理大数据集
  2. 内存管理:注意内存使用,对大文件使用流式处理
  3. 错误处理:实现完整的错误处理和恢复机制
  4. 日志记录:记录详细的执行日志
  5. 性能优化:使用索引、批量插入、并行处理
  6. 可重试性:确保ETL可以安全重试
  7. 测试:为每个ETL编写单元和集成测试
  8. 监控:监控执行时间、数据量、错误率

选择哪种方法取决于你的具体需求:

  • 小规模简单任务:手动实现
  • 中等规模:使用 PDO 自己封装
  • 生产级系统:使用成熟的 ETL 框架或工具

抱歉,评论功能暂时关闭!