flow-php/etl
Flow PHP ETL is a strongly typed, generator-powered ETL framework for efficient extract-transform-load pipelines in PHP. Process large datasets with a minimal memory footprint and plug into many adapters, extractors, and loaders for diverse sources.
Installation:
composer require flow-php/etl
Publish the config file:
php artisan vendor:publish --provider="Flow\ETL\ETLServiceProvider" --tag="config"
First Use Case: Define a simple ETL pipeline to extract data from a CSV, transform it, and load it into a database table.
use Flow\ETL\ETL;
use Flow\ETL\Extractor\CsvExtractor;
use Flow\ETL\Loader\DatabaseLoader;
use Flow\ETL\Transformer\ArrayTransformer;
$etl = app(ETL::class);
$etl->run(
new CsvExtractor('path/to/data.csv'),
new ArrayTransformer(fn($row) => ['name' => $row['first_name'], 'email' => $row['email']]),
new DatabaseLoader('users', connection: 'mysql')
);
Where to Look First:
config/etl.php for adapter configurations.vendor/flow-php/etl/src/ for source code examples (extractors, loaders, transformers).Pipeline Composition: Chain extractors, transformers, and loaders in a fluent manner.
$etl->run(
new ApiExtractor('https://api.example.com/users'),
new JsonPathTransformer('$.data[*].{name,email}'),
new DatabaseLoader('users')
);
Generator-Based Processing: Leverage generators for memory efficiency with large datasets.
$etl->run(
new DatabaseExtractor('SELECT * FROM large_table'),
new ChunkTransformer(1000), // Process in chunks
new S3Loader('bucket-name')
);
Adapter Customization: Extend built-in adapters or create custom ones.
class CustomDatabaseExtractor implements ExtractorInterface {
public function extract(): Generator {
yield ['id' => 1, 'name' => 'John'];
yield ['id' => 2, 'name' => 'Jane'];
}
}
Integration with Laravel Queues: Dispatch ETL jobs asynchronously.
use Flow\ETL\ETLJob;
ETLJob::dispatch(
new CsvExtractor('data.csv'),
new ArrayTransformer(...),
new DatabaseLoader('users')
)->onQueue('etl');
Error Handling and Retries: Use Laravel’s queue retries or custom middleware.
$etl->run(
new ApiExtractor('https://api.example.com/users'),
new RetryTransformer(3, fn($row) => $row['status'] === 'failed')
);
Data Migration:
LegacyDatabaseExtractor).ArrayTransformer).DatabaseLoader).Real-Time Analytics:
WebhookExtractor).AggregateTransformer).InfluxDBLoader).Automated Reporting:
MultiApiExtractor).ReportTransformer).CsvLoader, PdfLoader).Laravel Service Container: Bind adapters as Laravel services for dependency injection.
$this->app->bind(
ExtractorInterface::class,
fn($app) => new ApiExtractor('https://api.example.com')
);
Configuration: Use Laravel’s config system to manage adapter settings.
// config/etl.php
'adapters' => [
'api' => [
'base_url' => env('API_BASE_URL'),
'timeout' => 30,
],
],
Artisan Commands: Create CLI commands for running ETL pipelines.
class RunDataSync extends Command {
public function handle() {
$etl = app(ETL::class);
$etl->run(new DataSyncPipeline());
}
}
Event Listeners: Emit events for job lifecycle management.
// Listen for ETL job started/failed events
event(new EtlJobStarted($job));
Testing: Use Laravel’s testing tools to mock adapters and pipelines.
$this->mock(ExtractorInterface::class, function ($mock) {
$mock->shouldReceive('extract')->andReturn([['name' => 'Test']]);
});
Memory Leaks:
Connection facade for database operations to manage connections.Type Safety:
assert statements.
$transformer = new ArrayTransformer(fn(array $row): array => [
'name' => $row['first_name'] ?? '',
'email' => assert($row['email'] ?? null, 'string'),
]);
Queue Blocking:
chunk() method for database operations.Adapter Dependencies:
composer.json for unused dependencies. Use Laravel’s config/etl.php to disable unused adapters.Debugging Generators:
tap() to inspect intermediate data.
$etl->run(
new CsvExtractor('data.csv'),
new ArrayTransformer(fn($row) => tap($row, fn($r) => Log::debug('Row:', $r))),
new DatabaseLoader('users')
);
Idempotency:
UpsertDatabaseLoader).Performance Bottlenecks:
Leverage Laravel Facades:
Use Laravel’s facades (e.g., DB, Cache, Log) within adapters for consistency.
class DatabaseExtractor implements ExtractorInterface {
public function extract(): Generator {
foreach (DB::table('users')->cursor() as $user) {
yield $user;
}
}
}
Custom Middleware: Add middleware for cross-cutting concerns (e.g., logging, validation).
$etl->run(
new CsvExtractor('data.csv'),
new LoggingMiddleware(),
new ValidationTransformer(fn($row) => validator($row, ['email' => 'required|email'])),
new DatabaseLoader('users')
);
Environment-Specific Configs:
Use Laravel’s .env for environment-specific settings.
// .env
ETL_API_TIMEOUT=60
ETL_S3_BUCKET=my-bucket
Monitoring: Integrate with Laravel Horizon or custom dashboards to track job progress.
// Emit progress events
event(new EtlProgressEvent(50, 'Processing users...'));
Testing Strategies:
Extending Adapters: Create reusable adapter wrappers for common use cases.
class EloquentLoader implements LoaderInterface {
public function load(iterable $data): void {
foreach ($data as $item) {
User::create($item);
}
}
}
Documentation: Document complex pipelines with PHPDoc or
How can I help you explore Laravel packages today?