MidПрактика28 min

Data Pipeline: ETL и ELT

Различия ETL и ELT, инструменты, data warehouses, построение ETL pipeline на PHP

Что такое Data Pipeline

Data pipeline -- система для перемещения и трансформации данных из источников (sources) в хранилища (destinations). Два основных подхода: ETL и ELT.

ETL: Extract, Transform, Load

Данные извлекаются, трансформируются в промежуточном слое, затем загружаются в хранилище.

┌──────────┐    ┌──────────────┐    ┌──────────────┐
│  Sources │───>│  Transform   │───>│ Data         │
│          │    │  (staging)   │    │ Warehouse    │
│ - DB     │    │              │    │              │
│ - API    │    │ - Clean      │    │ - Star       │
│ - Files  │    │ - Validate   │    │   schema     │
│ - Logs   │    │ - Aggregate  │    │ - Analytics  │
└──────────┘    └──────────────┘    └──────────────┘
   Extract        Transform             Load

ETL Pipeline

<?php

declare(strict_types=1);

/**
 * ETL Pipeline framework in PHP
 */

// Step 1: Extract interface
interface Extractor
{
    /** @return iterable<array> */
    public function extract(): iterable;
}

// Step 2: Transform interface
interface Transformer
{
    public function transform(array $record): ?array;
}

// Step 3: Load interface
interface Loader
{
    public function load(array $record): void;
    public function flush(): void;
}

/**
 * Pipeline orchestrator
 */
final class EtlPipeline
{
    /** @var Transformer[] */
    private array $transformers = [];

    public function __construct(
        private readonly Extractor $extractor,
        private readonly Loader $loader,
        private readonly PipelineLogger $logger,
    ) {}

    public function addTransformer(Transformer $transformer): self
    {
        $this->transformers[] = $transformer;
        return $this;
    }

    public function run(): PipelineResult
    {
        $stats = ['extracted' => 0, 'transformed' => 0, 'loaded' => 0, 'errors' => 0];

        foreach ($this->extractor->extract() as $record) {
            $stats['extracted']++;

            try {
                $transformed = $record;

                foreach ($this->transformers as $transformer) {
                    $transformed = $transformer->transform($transformed);
                    if ($transformed === null) {
                        break; // Record filtered out
                    }
                }

                if ($transformed !== null) {
                    $this->loader->load($transformed);
                    $stats['loaded']++;
                }

                $stats['transformed']++;
            } catch (\Throwable $e) {
                $stats['errors']++;
                $this->logger->error("ETL error at record {$stats['extracted']}", [
                    'error' => $e->getMessage(),
                    'record' => $record,
                ]);
            }

            if ($stats['extracted'] % 1000 === 0) {
                $this->logger->info("Progress: {$stats['extracted']} records processed");
            }
        }

        $this->loader->flush();

        return new PipelineResult($stats);
    }
}

final readonly class PipelineResult
{
    public function __construct(
        public array $stats,
    ) {}
}

interface PipelineLogger
{
    public function info(string $message, array $context = []): void;
    public function error(string $message, array $context = []): void;
}
### Пример: ETL из API в PostgreSQL
<?php

declare(strict_types=1);

/**
 * Extract: fetch orders from external REST API
 */
final class ApiExtractor implements Extractor
{
    private int $page = 1;
    private readonly \DateTimeImmutable $since;

    public function __construct(
        private readonly string $apiUrl,
        private readonly string $apiKey,
    ) {
        $this->since = new \DateTimeImmutable('-1 day');
    }

    public function extract(): iterable
    {
        do {
            $response = $this->fetchPage($this->page);
            $data = json_decode($response, true, 512, JSON_THROW_ON_ERROR);

            foreach ($data['items'] as $item) {
                yield $item;
            }

            $this->page++;
            $hasMore = $data['has_next_page'] ?? false;

            // Rate limiting
            usleep(100_000); // 100ms between requests
        } while ($hasMore);
    }

    private function fetchPage(int $page): string
    {
        $ch = curl_init();
        curl_setopt_array($ch, [
            CURLOPT_URL => "{$this->apiUrl}/orders?" . http_build_query([
                'since' => $this->since->format('c'),
                'page' => $page,
                'per_page' => 100,
            ]),
            CURLOPT_HTTPHEADER => ["Authorization: Bearer {$this->apiKey}"],
            CURLOPT_RETURNTRANSFER => true,
            CURLOPT_TIMEOUT => 30,
        ]);

        $result = curl_exec($ch);
        $httpCode = curl_getinfo($ch, CURLINFO_HTTP_CODE);
        curl_close($ch);

        if ($httpCode !== 200) {
            throw new \RuntimeException("API returned HTTP {$httpCode}");
        }

        return $result;
    }
}

/**
 * Transform: clean, validate, enrich data
 */
final class OrderTransformer implements Transformer
{
    public function transform(array $record): ?array
    {
        // Filter: skip test orders
        if (str_starts_with($record['email'] ?? '', 'test@')) {
            return null;
        }

        // Clean and validate
        $amount = (float) ($record['total_amount'] ?? 0);
        if ($amount <= 0) {
            return null;
        }

        // Transform: normalize and enrich
        return [
            'order_id' => $record['id'],
            'customer_email' => mb_strtolower(trim($record['email'])),
            'amount' => $amount,
            'currency' => strtoupper($record['currency'] ?? 'USD'),
            'status' => $this->normalizeStatus($record['status']),
            'category' => $record['items'][0]['category'] ?? 'unknown',
            'item_count' => count($record['items'] ?? []),
            'created_date' => (new \DateTimeImmutable($record['created_at']))->format('Y-m-d'),
            'created_at' => $record['created_at'],
        ];
    }

    private function normalizeStatus(string $status): string
    {
        return match (strtolower($status)) {
            'paid', 'completed', 'fulfilled' => 'completed',
            'pending', 'processing' => 'pending',
            'refunded', 'returned' => 'refunded',
            'cancelled', 'canceled' => 'cancelled',
            default => 'unknown',
        };
    }
}

/**
 * Load: batch insert into PostgreSQL
 */
final class PostgresLoader implements Loader
{
    private array $buffer = [];
    private readonly int $batchSize;

    public function __construct(
        private readonly \PDO $db,
        int $batchSize = 500,
    ) {
        $this->batchSize = $batchSize;
    }

    public function load(array $record): void
    {
        $this->buffer[] = $record;

        if (count($this->buffer) >= $this->batchSize) {
            $this->flushBatch();
        }
    }

    public function flush(): void
    {
        if (!empty($this->buffer)) {
            $this->flushBatch();
        }
    }

    private function flushBatch(): void
    {
        if (empty($this->buffer)) {
            return;
        }

        $columns = array_keys($this->buffer[0]);
        $placeholders = [];
        $values = [];

        foreach ($this->buffer as $i => $record) {
            $rowPlaceholders = [];
            foreach ($columns as $col) {
                $key = "{$col}_{$i}";
                $rowPlaceholders[] = ":{$key}";
                $values[$key] = $record[$col];
            }
            $placeholders[] = '(' . implode(', ', $rowPlaceholders) . ')';
        }

        $columnList = implode(', ', $columns);
        $sql = <<<SQL
            INSERT INTO orders_warehouse ({$columnList})
            VALUES {$placeholderStr}
            ON CONFLICT (order_id) DO UPDATE SET
                status = EXCLUDED.status,
                amount = EXCLUDED.amount
        SQL;
        $placeholderStr = implode(', ', $placeholders);

        $stmt = $this->db->prepare(
            "INSERT INTO orders_warehouse ({$columnList}) VALUES {$placeholderStr}
             ON CONFLICT (order_id) DO UPDATE SET status = EXCLUDED.status, amount = EXCLUDED.amount"
        );
        $stmt->execute($values);

        $this->buffer = [];
    }
}
### Запуск ETL pipeline
<?php

declare(strict_types=1);

// Assemble and run the pipeline
$pipeline = new EtlPipeline(
    extractor: new ApiExtractor('https://api.shop.com/v1', $apiKey),
    loader: new PostgresLoader($pdo, batchSize: 500),
    logger: new StdoutLogger(),
);

$pipeline->addTransformer(new OrderTransformer());

$result = $pipeline->run();

echo "ETL completed: " . json_encode($result->stats);
// ETL completed: {"extracted":5432,"transformed":5100,"loaded":5100,"errors":12}
## ELT: Extract, Load, Transform

В ELT данные сначала загружаются в хранилище "как есть", затем трансформируются внутри хранилища средствами SQL.

┌──────────┐    ┌──────────────┐    ┌──────────────┐
│  Sources │───>│ Data         │───>│ Transformed  │
│          │    │ Warehouse    │    │ Tables       │
│          │    │ (raw zone)   │    │ (analytics)  │
└──────────┘    └──────────────┘    └──────────────┘
   Extract         Load              Transform
                                     (SQL inside DW)

Преимущества ELT

Критерий ETL ELT
Где трансформация Промежуточный сервер В хранилище (SQL)
Мощность Ограничена ETL-сервером Мощность DW (MPP)
Гибкость Логика фиксирована Данные сырые, можно менять
Скорость загрузки Медленнее (трансформация) Быстрее (прямая загрузка)
Хранение raw data Нет (только трансформированные) Да
Инструменты Informatica, Talend, PHP dbt, SQL, BigQuery

ELT -- загрузка сырых данных

<?php

declare(strict_types=1);

/**
 * ELT: load raw data first, transform later with SQL
 */
final class RawDataLoader
{
    public function __construct(
        private readonly \PDO $db,
    ) {}

    /**
     * Load raw JSON into staging table
     */
    public function loadRaw(string $source, iterable $records): int
    {
        $stmt = $this->db->prepare(<<<SQL
            INSERT INTO raw_events (source, raw_data, ingested_at)
            VALUES (:source, :raw_data, NOW())
        SQL);

        $count = 0;
        foreach ($records as $record) {
            $stmt->execute([
                'source' => $source,
                'raw_data' => json_encode($record, JSON_THROW_ON_ERROR),
            ]);
            $count++;
        }

        return $count;
    }

    /**
     * Transform inside database using SQL
     */
    public function transformWithSql(): void
    {
        // Create materialized view from raw data
        $this->db->exec(<<<SQL
            CREATE MATERIALIZED VIEW IF NOT EXISTS orders_analytics AS
            SELECT
                (raw_data->>'order_id')::TEXT as order_id,
                LOWER(raw_data->>'email') as customer_email,
                (raw_data->>'total_amount')::DECIMAL(10,2) as amount,
                UPPER(raw_data->>'currency') as currency,
                CASE LOWER(raw_data->>'status')
                    WHEN 'paid' THEN 'completed'
                    WHEN 'completed' THEN 'completed'
                    WHEN 'pending' THEN 'pending'
                    WHEN 'refunded' THEN 'refunded'
                    ELSE 'unknown'
                END as status,
                (raw_data->>'created_at')::TIMESTAMPTZ as created_at,
                DATE_TRUNC('day', (raw_data->>'created_at')::TIMESTAMPTZ) as created_date
            FROM raw_events
            WHERE source = 'shop_api'
              AND (raw_data->>'total_amount')::DECIMAL > 0
              AND raw_data->>'email' NOT LIKE 'test@%'
        SQL);

        // Refresh materialized view
        $this->db->exec('REFRESH MATERIALIZED VIEW CONCURRENTLY orders_analytics');
    }
}
## Инструменты Data Pipeline

Категории инструментов

Категория Инструменты Назначение
Orchestration Apache Airflow, Dagster, Prefect Планирование и мониторинг
ELT Transform dbt SQL-трансформации
Data Integration Fivetran, Airbyte, Stitch Коннекторы к источникам
Stream Processing Kafka, Flink, Spark Потоковая обработка
Data Warehouse BigQuery, Snowflake, Redshift Хранение и аналитика

PHP как ETL: когда это оправдано

  • Небольшие и средние объёмы данных (до миллионов записей)
  • Интеграция с существующим PHP-стеком
  • Кастомная бизнес-логика трансформации
  • Нет бюджета на специализированные ETL-инструменты
<?php

declare(strict_types=1);

/**
 * Scheduled ETL job runner with monitoring
 */
final class ScheduledEtlRunner
{
    /** @var array<string, EtlPipeline> */
    private array $pipelines = [];

    public function __construct(
        private readonly PipelineLogger $logger,
        private readonly MetricsClient $metrics,
    ) {}

    public function register(string $name, EtlPipeline $pipeline): void
    {
        $this->pipelines[$name] = $pipeline;
    }

    /**
     * Run specific pipeline with monitoring
     */
    public function execute(string $name): PipelineResult
    {
        if (!isset($this->pipelines[$name])) {
            throw new \InvalidArgumentException("Pipeline '{$name}' not found");
        }

        $startTime = microtime(true);
        $this->logger->info("Starting pipeline: {$name}");
        $this->metrics->gauge("etl.{$name}.running", 1);

        try {
            $result = $this->pipelines[$name]->run();

            $duration = microtime(true) - $startTime;
            $this->metrics->timing("etl.{$name}.duration", $duration);
            $this->metrics->counter("etl.{$name}.records", $result->stats['loaded']);
            $this->metrics->counter("etl.{$name}.errors", $result->stats['errors']);

            $this->logger->info("Pipeline {$name} completed", [
                'duration' => round($duration, 2),
                'stats' => $result->stats,
            ]);

            return $result;
        } catch (\Throwable $e) {
            $this->metrics->counter("etl.{$name}.failures", 1);
            $this->logger->error("Pipeline {$name} failed: {$e->getMessage()}");
            throw $e;
        } finally {
            $this->metrics->gauge("etl.{$name}.running", 0);
        }
    }
}
## Data Quality

Качество данных -- ключевой аспект любого pipeline:

<?php

declare(strict_types=1);

/**
 * Data quality checks for ETL pipeline
 */
final class DataQualityChecker
{
    /** @var array<string, callable> */
    private array $checks = [];

    public function addCheck(string $name, callable $check): self
    {
        $this->checks[$name] = $check;
        return $this;
    }

    /**
     * Run all checks after pipeline completion
     */
    public function validate(\PDO $db): array
    {
        $results = [];

        foreach ($this->checks as $name => $check) {
            try {
                $passed = $check($db);
                $results[$name] = ['passed' => $passed, 'error' => null];
            } catch (\Throwable $e) {
                $results[$name] = ['passed' => false, 'error' => $e->getMessage()];
            }
        }

        return $results;
    }
}

// Usage
$checker = new DataQualityChecker();

$checker->addCheck('no_nulls_in_amount', function (\PDO $db): bool {
    $count = $db->query("SELECT COUNT(*) FROM orders_warehouse WHERE amount IS NULL")->fetchColumn();
    return $count === 0;
});

$checker->addCheck('no_future_dates', function (\PDO $db): bool {
    $count = $db->query("SELECT COUNT(*) FROM orders_warehouse WHERE created_at > NOW()")->fetchColumn();
    return $count === 0;
});

$checker->addCheck('row_count_reasonable', function (\PDO $db): bool {
    $count = (int) $db->query("SELECT COUNT(*) FROM orders_warehouse WHERE created_date = CURRENT_DATE")->fetchColumn();
    return $count > 0 && $count < 1_000_000; // Expect between 1 and 1M orders/day
});

$results = $checker->validate($pdo);
> **Best Practice:** всегда запускайте проверки качества после ETL/ELT pipeline. Автоматизируйте алерты при провале проверок.

Итоги

  • ETL -- трансформация до загрузки, подходит для структурированных данных и ограниченных хранилищ
  • ELT -- загрузка сырых данных, трансформация в хранилище средствами SQL, подходит для cloud DW
  • PHP подходит для ETL pipeline малого и среднего масштаба
  • Data quality checks обязательны для любого pipeline
  • Для больших объёмов рассмотрите специализированные инструменты (Airflow, dbt, Fivetran)