<?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;
}
package etl
import (
"context"
"fmt"
"iter"
"log/slog"
)
// Extractor yields records from a data source.
type Extractor interface {
Extract(ctx context.Context) iter.Seq2[map[string]any, error]
}
// Transformer converts a record; returns nil to filter it out.
type Transformer interface {
Transform(record map[string]any) (map[string]any, error)
}
// Loader writes records to the destination.
type Loader interface {
Load(ctx context.Context, record map[string]any) error
Flush(ctx context.Context) error
}
// PipelineResult holds ETL run statistics.
type PipelineResult struct {
Extracted int `json:"extracted"`
Transformed int `json:"transformed"`
Loaded int `json:"loaded"`
Errors int `json:"errors"`
}
// Pipeline orchestrates extract, transform, load steps.
type Pipeline struct {
extractor Extractor
transformers []Transformer
loader Loader
logger *slog.Logger
}
func NewPipeline(e Extractor, l Loader, logger *slog.Logger) *Pipeline {
return &Pipeline{extractor: e, loader: l, logger: logger}
}
func (p *Pipeline) AddTransformer(t Transformer) *Pipeline {
p.transformers = append(p.transformers, t)
return p
}
// Run executes the ETL pipeline.
func (p *Pipeline) Run(ctx context.Context) (*PipelineResult, error) {
stats := &PipelineResult{}
for record, err := range p.extractor.Extract(ctx) {
if err != nil {
stats.Errors++
p.logger.Error("extract error", "error", err)
continue
}
stats.Extracted++
transformed := record
skip := false
for _, t := range p.transformers {
transformed, err = t.Transform(transformed)
if err != nil {
stats.Errors++
p.logger.Error("transform error", "error", err, "record", stats.Extracted)
skip = true
break
}
if transformed == nil {
skip = true
break
}
}
if !skip {
if err := p.loader.Load(ctx, transformed); err != nil {
stats.Errors++
p.logger.Error("load error", "error", err)
} else {
stats.Loaded++
}
}
stats.Transformed++
if stats.Extracted%1000 == 0 {
p.logger.Info("progress", "extracted", stats.Extracted)
}
}
if err := p.loader.Flush(ctx); err != nil {
return stats, fmt.Errorf("flush loader: %w", err)
}
return stats, nil
}
using Microsoft.Extensions.Logging;
namespace Etl;
// Extractor yields records from a data source.
public interface IExtractor
{
IAsyncEnumerable<Dictionary<string, object?>> ExtractAsync(CancellationToken ct = default);
}
// Transformer converts a record; returns null to filter it out.
public interface ITransformer
{
Dictionary<string, object?>? Transform(Dictionary<string, object?> record);
}
// Loader writes records to the destination.
public interface ILoader
{
Task LoadAsync(Dictionary<string, object?> record, CancellationToken ct = default);
Task FlushAsync(CancellationToken ct = default);
}
// PipelineResult holds ETL run statistics.
public sealed record PipelineResult
{
public int Extracted { get; set; }
public int Transformed { get; set; }
public int Loaded { get; set; }
public int Errors { get; set; }
}
// Pipeline orchestrates extract, transform, load steps.
public sealed class EtlPipeline(IExtractor extractor, ILoader loader, ILogger<EtlPipeline> logger)
{
private readonly List<ITransformer> _transformers = [];
public EtlPipeline AddTransformer(ITransformer transformer)
{
_transformers.Add(transformer);
return this;
}
public async Task<PipelineResult> RunAsync(CancellationToken ct = default)
{
var stats = new PipelineResult();
await foreach (var record in extractor.ExtractAsync(ct))
{
stats.Extracted++;
try
{
Dictionary<string, object?>? transformed = record;
foreach (var transformer in _transformers)
{
transformed = transformer.Transform(transformed);
if (transformed is null)
{
break; // Record filtered out
}
}
if (transformed is not null)
{
await loader.LoadAsync(transformed, ct);
stats.Loaded++;
}
stats.Transformed++;
}
catch (Exception ex)
{
stats.Errors++;
logger.LogError(ex, "ETL error at record {Record}", stats.Extracted);
}
if (stats.Extracted % 1000 == 0)
{
logger.LogInformation("Progress: {Extracted} records processed", stats.Extracted);
}
}
await loader.FlushAsync(ct);
return stats;
}
}
from __future__ import annotations
import logging
from abc import ABC, abstractmethod
from collections.abc import AsyncIterator
from dataclasses import dataclass
from typing import Any
logger = logging.getLogger(__name__)
Record = dict[str, Any]
class Extractor(ABC):
"""Yields records from a data source."""
@abstractmethod
def extract(self) -> AsyncIterator[Record]: ...
class Transformer(ABC):
"""Converts a record; returns None to filter it out."""
@abstractmethod
def transform(self, record: Record) -> Record | None: ...
class Loader(ABC):
"""Writes records to the destination."""
@abstractmethod
async def load(self, record: Record) -> None: ...
@abstractmethod
async def flush(self) -> None: ...
@dataclass
class PipelineResult:
"""Holds ETL run statistics."""
extracted: int = 0
transformed: int = 0
loaded: int = 0
errors: int = 0
class EtlPipeline:
"""Orchestrates extract, transform, load steps."""
def __init__(self, extractor: Extractor, loader: Loader) -> None:
self._extractor = extractor
self._loader = loader
self._transformers: list[Transformer] = []
def add_transformer(self, transformer: Transformer) -> EtlPipeline:
self._transformers.append(transformer)
return self
async def run(self) -> PipelineResult:
stats = PipelineResult()
async for record in self._extractor.extract():
stats.extracted += 1
try:
transformed: Record | None = record
for transformer in self._transformers:
transformed = transformer.transform(transformed)
if transformed is None:
break # Record filtered out
if transformed is not None:
await self._loader.load(transformed)
stats.loaded += 1
stats.transformed += 1
except Exception:
stats.errors += 1
logger.exception("ETL error at record %d", stats.extracted)
if stats.extracted % 1000 == 0:
logger.info("Progress: %d records processed", stats.extracted)
await self._loader.flush()
return stats
### Пример: 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 = [];
}
}
package etl
import (
"context"
"database/sql"
"encoding/json"
"fmt"
"iter"
"net/http"
"strings"
"time"
)
// ApiExtractor fetches orders from an external REST API.
type ApiExtractor struct {
apiURL string
apiKey string
since time.Time
}
func NewApiExtractor(apiURL, apiKey string) *ApiExtractor {
return &ApiExtractor{apiURL: apiURL, apiKey: apiKey, since: time.Now().Add(-24 * time.Hour)}
}
func (e *ApiExtractor) Extract(ctx context.Context) iter.Seq2[map[string]any, error] {
return func(yield func(map[string]any, error) bool) {
page := 1
for {
url := fmt.Sprintf("%s/orders?since=%s&page=%d&per_page=100",
e.apiURL, e.since.Format(time.RFC3339), page)
req, _ := http.NewRequestWithContext(ctx, "GET", url, nil)
req.Header.Set("Authorization", "Bearer "+e.apiKey)
resp, err := http.DefaultClient.Do(req)
if err != nil {
yield(nil, fmt.Errorf("fetch page %d: %w", page, err))
return
}
var data struct {
Items []map[string]any `json:"items"`
HasNextPage bool `json:"has_next_page"`
}
json.NewDecoder(resp.Body).Decode(&data)
resp.Body.Close()
for _, item := range data.Items {
if !yield(item, nil) {
return
}
}
if !data.HasNextPage {
return
}
page++
time.Sleep(100 * time.Millisecond) // Rate limiting
}
}
}
// OrderTransformer cleans, validates, and enriches order data.
type OrderTransformer struct{}
func (t *OrderTransformer) Transform(record map[string]any) (map[string]any, error) {
email, _ := record["email"].(string)
if strings.HasPrefix(email, "test@") {
return nil, nil // Filter test orders
}
amount, _ := record["total_amount"].(float64)
if amount <= 0 {
return nil, nil
}
currency, _ := record["currency"].(string)
if currency == "" {
currency = "USD"
}
status, _ := record["status"].(string)
category := "unknown"
if items, ok := record["items"].([]any); ok && len(items) > 0 {
if first, ok := items[0].(map[string]any); ok {
if c, ok := first["category"].(string); ok {
category = c
}
}
}
return map[string]any{
"order_id": record["id"],
"customer_email": strings.ToLower(strings.TrimSpace(email)),
"amount": amount,
"currency": strings.ToUpper(currency),
"status": normalizeStatus(status),
"category": category,
"created_at": record["created_at"],
}, nil
}
func normalizeStatus(status string) string {
switch strings.ToLower(status) {
case "paid", "completed", "fulfilled":
return "completed"
case "pending", "processing":
return "pending"
case "refunded", "returned":
return "refunded"
case "cancelled", "canceled":
return "cancelled"
default:
return "unknown"
}
}
// PostgresLoader batch-inserts records into PostgreSQL.
type PostgresLoader struct {
db *sql.DB
buffer []map[string]any
batchSize int
}
func NewPostgresLoader(db *sql.DB, batchSize int) *PostgresLoader {
return &PostgresLoader{db: db, batchSize: batchSize}
}
func (l *PostgresLoader) Load(ctx context.Context, record map[string]any) error {
l.buffer = append(l.buffer, record)
if len(l.buffer) >= l.batchSize {
return l.flushBatch(ctx)
}
return nil
}
func (l *PostgresLoader) Flush(ctx context.Context) error {
if len(l.buffer) > 0 {
return l.flushBatch(ctx)
}
return nil
}
func (l *PostgresLoader) flushBatch(ctx context.Context) error {
if len(l.buffer) == 0 {
return nil
}
tx, err := l.db.BeginTx(ctx, nil)
if err != nil {
return fmt.Errorf("begin tx: %w", err)
}
defer tx.Rollback()
stmt, err := tx.PrepareContext(ctx,
`INSERT INTO orders_warehouse (order_id, customer_email, amount, currency, status, category, created_at)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (order_id) DO UPDATE SET status = EXCLUDED.status, amount = EXCLUDED.amount`)
if err != nil {
return fmt.Errorf("prepare: %w", err)
}
defer stmt.Close()
for _, r := range l.buffer {
if _, err := stmt.ExecContext(ctx,
r["order_id"], r["customer_email"], r["amount"],
r["currency"], r["status"], r["category"], r["created_at"],
); err != nil {
return fmt.Errorf("exec insert: %w", err)
}
}
l.buffer = l.buffer[:0]
return tx.Commit()
}
using System.Net.Http.Json;
using System.Runtime.CompilerServices;
using System.Text.Json;
using Npgsql;
namespace Etl;
// ApiExtractor fetches orders from an external REST API.
public sealed class ApiExtractor(HttpClient http, string apiUrl, string apiKey) : IExtractor
{
private readonly DateTimeOffset _since = DateTimeOffset.UtcNow.AddDays(-1);
public async IAsyncEnumerable<Dictionary<string, object?>> ExtractAsync(
[EnumeratorCancellation] CancellationToken ct = default)
{
var page = 1;
bool hasMore;
do
{
var url = $"{apiUrl}/orders?since={Uri.EscapeDataString(_since.ToString("o"))}" +
$"&page={page}&per_page=100";
using var request = new HttpRequestMessage(HttpMethod.Get, url);
request.Headers.Add("Authorization", $"Bearer {apiKey}");
using var response = await http.SendAsync(request, ct);
response.EnsureSuccessStatusCode();
var data = await response.Content.ReadFromJsonAsync<ApiPage>(ct)
?? throw new InvalidOperationException("Empty API response");
foreach (var item in data.Items)
{
yield return item;
}
page++;
hasMore = data.HasNextPage;
await Task.Delay(TimeSpan.FromMilliseconds(100), ct); // Rate limiting
} while (hasMore);
}
private sealed record ApiPage(
List<Dictionary<string, object?>> Items,
bool HasNextPage);
}
// OrderTransformer cleans, validates, and enriches order data.
public sealed class OrderTransformer : ITransformer
{
public Dictionary<string, object?>? Transform(Dictionary<string, object?> record)
{
var email = record.GetValueOrDefault("email") as string ?? "";
if (email.StartsWith("test@", StringComparison.OrdinalIgnoreCase))
{
return null; // Filter test orders
}
var amount = Convert.ToDecimal(record.GetValueOrDefault("total_amount") ?? 0m);
if (amount <= 0)
{
return null;
}
var currency = record.GetValueOrDefault("currency") as string ?? "USD";
var status = record.GetValueOrDefault("status") as string ?? "";
var category = "unknown";
if (record.GetValueOrDefault("items") is JsonElement { ValueKind: JsonValueKind.Array } items
&& items.GetArrayLength() > 0
&& items[0].TryGetProperty("category", out var cat))
{
category = cat.GetString() ?? "unknown";
}
return new Dictionary<string, object?>
{
["order_id"] = record.GetValueOrDefault("id"),
["customer_email"] = email.Trim().ToLowerInvariant(),
["amount"] = amount,
["currency"] = currency.ToUpperInvariant(),
["status"] = NormalizeStatus(status),
["category"] = category,
["created_at"] = record.GetValueOrDefault("created_at"),
};
}
private static string NormalizeStatus(string status) => status.ToLowerInvariant() switch
{
"paid" or "completed" or "fulfilled" => "completed",
"pending" or "processing" => "pending",
"refunded" or "returned" => "refunded",
"cancelled" or "canceled" => "cancelled",
_ => "unknown",
};
}
// PostgresLoader batch-inserts records into PostgreSQL.
public sealed class PostgresLoader(NpgsqlDataSource dataSource, int batchSize = 500) : ILoader
{
private readonly List<Dictionary<string, object?>> _buffer = [];
public async Task LoadAsync(Dictionary<string, object?> record, CancellationToken ct = default)
{
_buffer.Add(record);
if (_buffer.Count >= batchSize)
{
await FlushBatchAsync(ct);
}
}
public async Task FlushAsync(CancellationToken ct = default)
{
if (_buffer.Count > 0)
{
await FlushBatchAsync(ct);
}
}
private async Task FlushBatchAsync(CancellationToken ct)
{
await using var conn = await dataSource.OpenConnectionAsync(ct);
await using var tx = await conn.BeginTransactionAsync(ct);
foreach (var r in _buffer)
{
await using var cmd = new NpgsqlCommand(
"""
INSERT INTO orders_warehouse
(order_id, customer_email, amount, currency, status, category, created_at)
VALUES (@order_id, @email, @amount, @currency, @status, @category, @created_at)
ON CONFLICT (order_id) DO UPDATE SET
status = EXCLUDED.status,
amount = EXCLUDED.amount
""", conn, tx);
cmd.Parameters.AddWithValue("order_id", r["order_id"] ?? DBNull.Value);
cmd.Parameters.AddWithValue("email", r["customer_email"] ?? DBNull.Value);
cmd.Parameters.AddWithValue("amount", r["amount"] ?? DBNull.Value);
cmd.Parameters.AddWithValue("currency", r["currency"] ?? DBNull.Value);
cmd.Parameters.AddWithValue("status", r["status"] ?? DBNull.Value);
cmd.Parameters.AddWithValue("category", r["category"] ?? DBNull.Value);
cmd.Parameters.AddWithValue("created_at", r["created_at"] ?? DBNull.Value);
await cmd.ExecuteNonQueryAsync(ct);
}
await tx.CommitAsync(ct);
_buffer.Clear();
}
}
from __future__ import annotations
import asyncio
from collections.abc import AsyncIterator
from datetime import UTC, datetime, timedelta
from decimal import Decimal
from typing import Any
import asyncpg
import httpx
Record = dict[str, Any]
class ApiExtractor(Extractor):
"""Fetches orders from an external REST API."""
def __init__(self, client: httpx.AsyncClient, api_url: str, api_key: str) -> None:
self._client = client
self._api_url = api_url
self._api_key = api_key
self._since = datetime.now(UTC) - timedelta(days=1)
async def extract(self) -> AsyncIterator[Record]:
page = 1
while True:
response = await self._client.get(
f"{self._api_url}/orders",
params={
"since": self._since.isoformat(),
"page": page,
"per_page": 100,
},
headers={"Authorization": f"Bearer {self._api_key}"},
timeout=30.0,
)
response.raise_for_status()
data = response.json()
for item in data["items"]:
yield item
if not data.get("has_next_page", False):
return
page += 1
await asyncio.sleep(0.1) # Rate limiting
class OrderTransformer(Transformer):
"""Cleans, validates, and enriches order data."""
_STATUS_MAP = {
"paid": "completed",
"completed": "completed",
"fulfilled": "completed",
"pending": "pending",
"processing": "pending",
"refunded": "refunded",
"returned": "refunded",
"cancelled": "cancelled",
"canceled": "cancelled",
}
def transform(self, record: Record) -> Record | None:
email = record.get("email") or ""
if email.startswith("test@"):
return None # Filter test orders
amount = Decimal(str(record.get("total_amount") or 0))
if amount <= 0:
return None
items = record.get("items") or []
category = items[0].get("category", "unknown") if items else "unknown"
return {
"order_id": record["id"],
"customer_email": email.strip().lower(),
"amount": amount,
"currency": (record.get("currency") or "USD").upper(),
"status": self._normalize_status(record.get("status") or ""),
"category": category,
"created_at": record["created_at"],
}
def _normalize_status(self, status: str) -> str:
return self._STATUS_MAP.get(status.lower(), "unknown")
class PostgresLoader(Loader):
"""Batch-inserts records into PostgreSQL."""
_COLUMNS = (
"order_id", "customer_email", "amount",
"currency", "status", "category", "created_at",
)
def __init__(self, pool: asyncpg.Pool, batch_size: int = 500) -> None:
self._pool = pool
self._batch_size = batch_size
self._buffer: list[Record] = []
async def load(self, record: Record) -> None:
self._buffer.append(record)
if len(self._buffer) >= self._batch_size:
await self._flush_batch()
async def flush(self) -> None:
if self._buffer:
await self._flush_batch()
async def _flush_batch(self) -> None:
rows = [tuple(r[col] for col in self._COLUMNS) for r in self._buffer]
async with self._pool.acquire() as conn, conn.transaction():
await conn.executemany(
"""
INSERT INTO orders_warehouse
(order_id, customer_email, amount, currency, status, category, created_at)
VALUES ($1, $2, $3, $4, $5, $6, $7)
ON CONFLICT (order_id) DO UPDATE SET
status = EXCLUDED.status,
amount = EXCLUDED.amount
""",
rows,
)
self._buffer.clear()
### Запуск 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}
using System.Text.Json;
using Microsoft.Extensions.Logging;
using Npgsql;
// Assemble and run the pipeline
using var loggerFactory = LoggerFactory.Create(b => b.AddSimpleConsole());
using var http = new HttpClient();
await using var dataSource = NpgsqlDataSource.Create(
Environment.GetEnvironmentVariable("DATABASE_URL")!);
var pipeline = new Etl.EtlPipeline(
new Etl.ApiExtractor(http, "https://api.shop.com/v1",
Environment.GetEnvironmentVariable("API_KEY")!),
new Etl.PostgresLoader(dataSource, batchSize: 500),
loggerFactory.CreateLogger<Etl.EtlPipeline>());
pipeline.AddTransformer(new Etl.OrderTransformer());
var result = await pipeline.RunAsync();
Console.WriteLine($"ETL completed: {JsonSerializer.Serialize(result)}");
// ETL completed: {"Extracted":5432,"Transformed":5100,"Loaded":5100,"Errors":12}
import asyncio
import dataclasses
import json
import logging
import os
import asyncpg
import httpx
async def main() -> None:
logging.basicConfig(level=logging.INFO)
# Assemble and run the pipeline
async with httpx.AsyncClient() as client:
pool = await asyncpg.create_pool(os.environ["DATABASE_URL"])
try:
pipeline = EtlPipeline(
extractor=ApiExtractor(client, "https://api.shop.com/v1", os.environ["API_KEY"]),
loader=PostgresLoader(pool, batch_size=500),
)
pipeline.add_transformer(OrderTransformer())
result = await pipeline.run()
print(f"ETL completed: {json.dumps(dataclasses.asdict(result))}")
# ETL completed: {"extracted":5432,"transformed":5100,"loaded":5100,"errors":12}
finally:
await pool.close()
if __name__ == "__main__":
asyncio.run(main())
## ELT: Extract, Load, Transform
В ELT данные сначала загружаются в хранилище "как есть", затем трансформируются внутри хранилища средствами SQL.
<?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');
}
}
package elt
import (
"context"
"database/sql"
"encoding/json"
"fmt"
)
// RawDataLoader loads raw data first, transforms later with SQL.
type RawDataLoader struct {
db *sql.DB
}
func NewRawDataLoader(db *sql.DB) *RawDataLoader {
return &RawDataLoader{db: db}
}
// LoadRaw inserts raw JSON records into a staging table.
func (l *RawDataLoader) LoadRaw(ctx context.Context, source string, records []map[string]any) (int, error) {
stmt, err := l.db.PrepareContext(ctx,
`INSERT INTO raw_events (source, raw_data, ingested_at) VALUES ($1, $2, NOW())`)
if err != nil {
return 0, fmt.Errorf("prepare: %w", err)
}
defer stmt.Close()
count := 0
for _, record := range records {
data, _ := json.Marshal(record)
if _, err := stmt.ExecContext(ctx, source, string(data)); err != nil {
return count, fmt.Errorf("insert raw event: %w", err)
}
count++
}
return count, nil
}
// TransformWithSQL creates and refreshes materialized views using SQL.
func (l *RawDataLoader) TransformWithSQL(ctx context.Context) error {
_, err := l.db.ExecContext(ctx, `
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@%'`)
if err != nil {
return fmt.Errorf("create materialized view: %w", err)
}
_, err = l.db.ExecContext(ctx, `REFRESH MATERIALIZED VIEW CONCURRENTLY orders_analytics`)
return err
}
using System.Text.Json;
using Npgsql;
namespace Elt;
// RawDataLoader loads raw data first, transforms later with SQL.
public sealed class RawDataLoader(NpgsqlDataSource dataSource)
{
// LoadRaw inserts raw JSON records into a staging table.
public async Task<int> LoadRawAsync(
string source,
IEnumerable<Dictionary<string, object?>> records,
CancellationToken ct = default)
{
await using var conn = await dataSource.OpenConnectionAsync(ct);
var count = 0;
foreach (var record in records)
{
await using var cmd = new NpgsqlCommand(
"""
INSERT INTO raw_events (source, raw_data, ingested_at)
VALUES (@source, @raw_data::jsonb, NOW())
""", conn);
cmd.Parameters.AddWithValue("source", source);
cmd.Parameters.AddWithValue("raw_data", JsonSerializer.Serialize(record));
await cmd.ExecuteNonQueryAsync(ct);
count++;
}
return count;
}
// TransformWithSql creates and refreshes materialized views using SQL.
public async Task TransformWithSqlAsync(CancellationToken ct = default)
{
await using var conn = await dataSource.OpenConnectionAsync(ct);
await using (var create = new NpgsqlCommand(
"""
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@%'
""", conn))
{
await create.ExecuteNonQueryAsync(ct);
}
await using var refresh = new NpgsqlCommand(
"REFRESH MATERIALIZED VIEW CONCURRENTLY orders_analytics", conn);
await refresh.ExecuteNonQueryAsync(ct);
}
}
from __future__ import annotations
import json
from collections.abc import Iterable
from typing import Any
import asyncpg
Record = dict[str, Any]
class RawDataLoader:
"""ELT: load raw data first, transform later with SQL."""
def __init__(self, pool: asyncpg.Pool) -> None:
self._pool = pool
async def load_raw(self, source: str, records: Iterable[Record]) -> int:
"""Insert raw JSON records into a staging table."""
rows = [(source, json.dumps(record)) for record in records]
async with self._pool.acquire() as conn:
await conn.executemany(
"""
INSERT INTO raw_events (source, raw_data, ingested_at)
VALUES ($1, $2::jsonb, NOW())
""",
rows,
)
return len(rows)
async def transform_with_sql(self) -> None:
"""Transform inside the database using SQL."""
async with self._pool.acquire() as conn:
await conn.execute(
"""
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@%'
"""
)
await conn.execute(
"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);
}
}
}
package etl
import (
"context"
"fmt"
"log/slog"
"time"
)
// MetricsClient records pipeline metrics.
type MetricsClient interface {
Gauge(name string, value float64)
Timing(name string, duration time.Duration)
Counter(name string, value int)
}
// ScheduledRunner manages and executes named ETL pipelines with monitoring.
type ScheduledRunner struct {
pipelines map[string]*Pipeline
logger *slog.Logger
metrics MetricsClient
}
func NewScheduledRunner(logger *slog.Logger, metrics MetricsClient) *ScheduledRunner {
return &ScheduledRunner{
pipelines: make(map[string]*Pipeline),
logger: logger,
metrics: metrics,
}
}
func (r *ScheduledRunner) Register(name string, p *Pipeline) {
r.pipelines[name] = p
}
// Execute runs a named pipeline with monitoring.
func (r *ScheduledRunner) Execute(ctx context.Context, name string) (*PipelineResult, error) {
p, ok := r.pipelines[name]
if !ok {
return nil, fmt.Errorf("pipeline %q not found", name)
}
start := time.Now()
r.logger.Info("starting pipeline", "name", name)
r.metrics.Gauge("etl."+name+".running", 1)
defer r.metrics.Gauge("etl."+name+".running", 0)
result, err := p.Run(ctx)
duration := time.Since(start)
if err != nil {
r.metrics.Counter("etl."+name+".failures", 1)
r.logger.Error("pipeline failed", "name", name, "error", err)
return result, err
}
r.metrics.Timing("etl."+name+".duration", duration)
r.metrics.Counter("etl."+name+".records", result.Loaded)
r.metrics.Counter("etl."+name+".errors", result.Errors)
r.logger.Info("pipeline completed", "name", name,
"duration", duration.Round(time.Millisecond), "loaded", result.Loaded)
return result, nil
}
## 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);
package etl
import (
"context"
"database/sql"
"fmt"
)
// CheckResult holds the outcome of a single quality check.
type CheckResult struct {
Passed bool `json:"passed"`
Error string `json:"error,omitempty"`
}
// QualityCheck is a function that validates data quality.
type QualityCheck func(ctx context.Context, db *sql.DB) (bool, error)
// DataQualityChecker runs validation checks after pipeline completion.
type DataQualityChecker struct {
checks map[string]QualityCheck
}
func NewDataQualityChecker() *DataQualityChecker {
return &DataQualityChecker{checks: make(map[string]QualityCheck)}
}
func (c *DataQualityChecker) AddCheck(name string, check QualityCheck) {
c.checks[name] = check
}
// Validate runs all registered checks and returns results.
func (c *DataQualityChecker) Validate(ctx context.Context, db *sql.DB) map[string]CheckResult {
results := make(map[string]CheckResult, len(c.checks))
for name, check := range c.checks {
passed, err := check(ctx, db)
if err != nil {
results[name] = CheckResult{Passed: false, Error: err.Error()}
} else {
results[name] = CheckResult{Passed: passed}
}
}
return results
}
// Example checks
func NoNullAmounts(ctx context.Context, db *sql.DB) (bool, error) {
var count int
err := db.QueryRowContext(ctx,
"SELECT COUNT(*) FROM orders_warehouse WHERE amount IS NULL").Scan(&count)
return count == 0, err
}
func NoFutureDates(ctx context.Context, db *sql.DB) (bool, error) {
var count int
err := db.QueryRowContext(ctx,
"SELECT COUNT(*) FROM orders_warehouse WHERE created_at > NOW()").Scan(&count)
return count == 0, err
}
func ReasonableRowCount(ctx context.Context, db *sql.DB) (bool, error) {
var count int
err := db.QueryRowContext(ctx,
"SELECT COUNT(*) FROM orders_warehouse WHERE created_date = CURRENT_DATE").Scan(&count)
return count > 0 && count < 1_000_000, err
}
using Npgsql;
namespace Etl;
// CheckResult holds the outcome of a single quality check.
public sealed record CheckResult(bool Passed, string? Error = null);
// QualityCheck validates data quality.
public delegate Task<bool> QualityCheck(NpgsqlDataSource dataSource, CancellationToken ct);
// DataQualityChecker runs validation checks after pipeline completion.
public sealed class DataQualityChecker
{
private readonly Dictionary<string, QualityCheck> _checks = [];
public DataQualityChecker AddCheck(string name, QualityCheck check)
{
_checks[name] = check;
return this;
}
// Validate runs all registered checks and returns results.
public async Task<Dictionary<string, CheckResult>> ValidateAsync(
NpgsqlDataSource dataSource,
CancellationToken ct = default)
{
var results = new Dictionary<string, CheckResult>(_checks.Count);
foreach (var (name, check) in _checks)
{
try
{
results[name] = new CheckResult(await check(dataSource, ct));
}
catch (Exception ex)
{
results[name] = new CheckResult(false, ex.Message);
}
}
return results;
}
private static async Task<long> ScalarAsync(
NpgsqlDataSource dataSource, string sql, CancellationToken ct)
{
await using var cmd = dataSource.CreateCommand(sql);
return Convert.ToInt64(await cmd.ExecuteScalarAsync(ct));
}
// Example checks
public static async Task<bool> NoNullAmountsAsync(NpgsqlDataSource ds, CancellationToken ct) =>
await ScalarAsync(ds,
"SELECT COUNT(*) FROM orders_warehouse WHERE amount IS NULL", ct) == 0;
public static async Task<bool> NoFutureDatesAsync(NpgsqlDataSource ds, CancellationToken ct) =>
await ScalarAsync(ds,
"SELECT COUNT(*) FROM orders_warehouse WHERE created_at > NOW()", ct) == 0;
public static async Task<bool> ReasonableRowCountAsync(NpgsqlDataSource ds, CancellationToken ct)
{
var count = await ScalarAsync(ds,
"SELECT COUNT(*) FROM orders_warehouse WHERE created_date = CURRENT_DATE", ct);
return count is > 0 and < 1_000_000; // Expect between 1 and 1M orders/day
}
}
from __future__ import annotations
import asyncio
import os
from collections.abc import Awaitable, Callable
from dataclasses import dataclass
import asyncpg
QualityCheck = Callable[[asyncpg.Pool], Awaitable[bool]]
@dataclass(frozen=True)
class CheckResult:
"""Holds the outcome of a single quality check."""
passed: bool
error: str | None = None
class DataQualityChecker:
"""Runs validation checks after pipeline completion."""
def __init__(self) -> None:
self._checks: dict[str, QualityCheck] = {}
def add_check(self, name: str, check: QualityCheck) -> DataQualityChecker:
self._checks[name] = check
return self
async def validate(self, pool: asyncpg.Pool) -> dict[str, CheckResult]:
results: dict[str, CheckResult] = {}
for name, check in self._checks.items():
try:
results[name] = CheckResult(passed=await check(pool))
except Exception as exc:
results[name] = CheckResult(passed=False, error=str(exc))
return results
# Usage
async def no_nulls_in_amount(pool: asyncpg.Pool) -> bool:
count = await pool.fetchval(
"SELECT COUNT(*) FROM orders_warehouse WHERE amount IS NULL"
)
return count == 0
async def no_future_dates(pool: asyncpg.Pool) -> bool:
count = await pool.fetchval(
"SELECT COUNT(*) FROM orders_warehouse WHERE created_at > NOW()"
)
return count == 0
async def row_count_reasonable(pool: asyncpg.Pool) -> bool:
count = await pool.fetchval(
"SELECT COUNT(*) FROM orders_warehouse WHERE created_date = CURRENT_DATE"
)
return 0 < count < 1_000_000 # Expect between 1 and 1M orders/day
async def main() -> None:
checker = (
DataQualityChecker()
.add_check("no_nulls_in_amount", no_nulls_in_amount)
.add_check("no_future_dates", no_future_dates)
.add_check("row_count_reasonable", row_count_reasonable)
)
# Top-level await is only valid inside a coroutine, so the entry point
# has to be one — unlike a REPL or a Jupyter cell.
pool = await asyncpg.create_pool(os.environ["DATABASE_URL"])
try:
results = await checker.validate(pool)
finally:
await pool.close()
for result in results:
print(f"{result.name}: {'ok' if result.passed else 'FAILED'}")
asyncio.run(main())
> **Best Practice:** всегда запускайте проверки качества после ETL/ELT pipeline. Автоматизируйте алерты при провале проверок.
Итоги
ETL -- трансформация до загрузки, подходит для структурированных данных и ограниченных хранилищ
ELT -- загрузка сырых данных, трансформация в хранилище средствами SQL, подходит для cloud DW
PHP подходит для ETL pipeline малого и среднего масштаба
Data quality checks обязательны для любого pipeline
Для больших объёмов рассмотрите специализированные инструменты (Airflow, dbt, Fivetran)