Search by

robelsust / laravel-data-pipeline

robelsust

Streaming ETL data pipeline for Laravel with transactional checkpoints, batch validation, deduplication and quarantine recovery.

Package info

github.com/robelsust/laravel-data-pipeline

pkg:composer/robelsust/laravel-data-pipeline

Statistics

Installs: 1

Dependents: 0

Suggesters: 0

Stars: 0

Open Issues: 0

v1.0.1 2026-10-06 18:15 UTC

This package is auto-updated.

Last update: 2026-10-09 04:20:04 UTC


README

Laravel Data Pipeline Architecture

Build Status Latest Version PHP Version Laravel Version Software License

Laravel Data Pipeline is a streaming ETL orchestration package designed for high-throughput batch and stream processing in Laravel.

It eliminates fragile, bespoke loops and solves the common failure modes of large-scale data ingestion: out-of-memory crashes, N+1 validation queries, partial corruption, unhandled poison pills, and unrecoverable mid-stream job crashes.

Key Features

  • โšก Bulk Persistence: Multi-row inserts and upserts; throughput depends on your database, stages, and chunk size.
  • ๐Ÿ›ก๏ธ Bounded Streaming Execution: Synchronous processing buffers one chunk, a bounded deduplication window, and up to 5,000 returned failures. Queue dispatch currently stages all job payloads in memory.
  • ๐ŸŽฏ Batched Database Validation: BatchValidator pre-fetches unique and exists database constraints in a single chunk query instead of firing individual database queries for every row.
  • ๐Ÿ”’ Atomic Transactional Chunks: Each chunk executes in an isolated database transaction. If a chunk fails, it rolls back cleanly with zero phantom records.
  • โ˜ฃ๏ธ Poison Pill Quarantine: Corrupted records are automatically isolated into the pipeline_errors table with redacted raw payloads, record indices, and validation errors without stalling the rest of the stream.
  • ๐Ÿ”„ Surgical Zero-Loss Recovery (pipeline:retry-failed): Re-process only the quarantined failed records after updating business rules, using transactional row attribution and quarantine cleanup.
  • ๐Ÿ›‘ Mid-Stream Crash Recovery (pipeline:resume): If a worker crashes or gets killed midway, resume cleanly from the exact last committed chunk checkpoint.
  • ๐Ÿš€ Distributed Queue Scaling (DataPipeline::queue()): Automatically chunk and fan-out workloads across parallel queue workers with database lease locks.
  • ๐Ÿงช Zero-Write Simulation Engine (dryRun() & preview()): Inspect and validate transformed records without writing a single byte to the database.

Historical Benchmark Performance (Embedded Implementation)

These historical results describe the earlier application implementation, not a throughput guarantee for this package release. Run the benchmark against your own schema and infrastructure.

Benchmark Scenario Dataset Size Duration Throughput Peak Memory Result
High-Speed Bulk Insert 1,000,000 rows 19.45s 51,409 rows/sec 46 MB 100% Persisted
Collision Merge (UpsertSink) 1,000,000 rows (250k old + 750k new) 30.11s 33,206 rows/sec 48 MB 250k updated in-place, 750k inserted
Poison Pill Quarantine 500,000 rows (2,500 poison pills) 40.12s 12,463 rows/sec 44 MB 497.5k clean, 2.5k quarantined
Surgical Quarantine Retry 2,500 quarantined rows 0.65s 3,822 rows/sec 34 MB 100% recovered (0 loss)
Crash & Resume Checkpoint 500,000 rows (50% crash injection) 6.56s (resume) 38,110 rows/sec 42 MB Chunks 1โ€“50 skipped, 51โ€“100 completed

Requirements

  • Laravel 12.x with PHP 8.2+, or Laravel 13.x with PHP 8.3+
  • Laravel 11 is no longer supported: its security support ended March 12, 2026. See the framework support policy.
  • MySQL 8+, PostgreSQL 14+, or SQLite 3.35+
  • For Redis: the PhpRedis extension or composer require predis/predis in the host application

Installation

Install the package via Composer:

composer require robelsust/laravel-data-pipeline

Run the database migrations to create the pipeline tracking and quarantine tables:

php artisan migrate

(Optional) Publish the configuration file:

php artisan vendor:publish --tag=data-pipeline-config

Quick Start (30 Seconds)

use Robel\DataPipeline\Facades\DataPipeline;

$result = DataPipeline::from(storage_path('imports/users.csv'))
    ->chunk(2500)
    ->mapFields([
        'full_name'  => 'name',
        'user_email' => 'email',
    ])
    ->transform(fn (array $record) => array_merge($record, [
        'email' => strtolower($record['email']),
    ]))
    ->validate([
        'name'  => ['required', 'string', 'max:255'],
        'email' => ['required', 'email', 'unique:users,email'],
    ])
    ->allowPartialSuccess()
    ->track()
    ->upsert('users', uniqueBy: ['email'], update: ['name', 'updated_at'])
    ->run();

echo "Processed {$result->successful} records in {$result->durationInSeconds}s!";

Comprehensive Architecture & API Guide

1. Sources

DataPipeline accepts any data source that implements SourceInterface or native PHP/Laravel iterables:

// 1. Streaming CSV (fgetcsv generator stream)
DataPipeline::from(storage_path('data/catalog.csv'));

// 2. Eloquent Cursor (streaming database rows)
DataPipeline::from(User::query()->where('active', true));

// 3. LazyCollection / PHP Generator
DataPipeline::from(LazyCollection::make(function () {
    yield from fetchExternalApiRows();
}));

// 4. Memory-bounded Custom Source
DataPipeline::from(new CustomWebhookSource($feedUrl));

2. Processing Stages

Stages are decoupled, single-responsibility transformers executed on every record or chunk:

$pipeline
    // Rename/map column keys
    ->mapFields(['cust_name' => 'name', 'mail' => 'email'])

    // Inline row transformation
    ->transform(function (array $data): array {
        $data['name'] = ucwords(trim($data['name']));
        $data['account_type'] = $data['tier'] > 2 ? 'enterprise' : 'standard';
        return $data;
    })

    // Drop rows that don't match criteria
    ->filter(fn (array $data) => ! empty($data['email']))

    // Batch validation (Zero N+1 database queries)
    ->validate([
        'name'  => ['required', 'string'],
        'email' => ['required', 'email', 'unique:users,email'],
    ])

    // In-stream deduplication (drop repeated keys within chunk or across stream)
    ->deduplicateBy('email', action: 'drop', store: 'redis');

3. Sinks (Persistence)

Persist processed records cleanly into any database or external service:

// High-speed multi-row bulk insert via QueryBuilder
$pipeline->insert('users');

// Bulk Upsert with conflict resolution (ON DUPLICATE KEY UPDATE)
$pipeline->upsert('users', uniqueBy: ['email'], update: ['name', 'updated_at']);

// Eloquent Model Sink (triggers model events, mutators, and custom casts)
$pipeline->to(new ModelSink(User::class));

// Custom Destination Callback
$pipeline->persist(function (RecordChunk $chunk) {
    Http::post('https://api.warehouse.internal/sync', $chunk->validPayloads())->throw();
});

4. Resilience & Error Handling

$pipeline
    // Isolate bad rows into quarantine while letting clean rows commit
    ->allowPartialSuccess(true)

    // OR: strict fail-all mode (chunk rollback on any error)
    ->failIfAny()

    // Redact sensitive attributes from quarantine payloads
    ->redactAttributes(['password', 'card_number', 'ssn'])

    // Real-time progress monitoring
    ->onProgress(function (int $processed, ?int $total, int $success, int $failed) {
        echo "Processed {$processed}/{$total} records (Success: {$success}, Failed: {$failed})\n";
    });

Artisan CLI Suite

The package includes production management commands out of the box. Resume and retry require a named definition registered in a service provider, or an explicit --class / --definition option. Register definitions in every worker process.

DataPipeline::register('users-import-v1', fn () => DataPipeline::from(storage_path('imports/users.csv'))
    ->name('users-import-v1')->as('users-import-v1')->track()->insert('users'));

1. List Recent Pipeline Runs

php artisan pipeline:list --limit=20 --status=COMPLETED

2. Inspect Run Diagnostics & Quarantined Errors

php artisan pipeline:show <run-uuid>

3. Surgically Retry Quarantined Records

php artisan pipeline:retry-failed <run-uuid>

4. Resume an Interrupted or Crashed Run

php artisan pipeline:resume <run-uuid>

5. High-Volume Ingestion Benchmark

php artisan pipeline:benchmark --records=1000000 --chunk-size=5000 --track --upsert

6. Resilience Stress Testing (Poison Pills & Crashes)

php artisan pipeline:stress-test all --records=500000 --chunk-size=5000

Production configuration and guarantees

The facade resolves PipelineManager; synchronous runs, queued jobs, and quarantine retries use the same ChunkExecutor. Configure data-pipeline.connection before migrations if tracking uses a non-default database.

  • Keep sink writes and tracking on the same connection for atomic commits. allowCrossConnectionTransactions() permits separate commits and cannot provide distributed atomicity. HTTP calls and other external effects need destination-side idempotency.
  • Use a shared lock-capable cache such as Redis for multiple workers. Tracked deduplication ownership is also journaled in pipeline_deduplication_claims inside the sink/checkpoint transaction, so committed keys survive cache loss. Untracked in-memory deduplication is a bounded sliding window, and untracked cache entries expire.
  • Match worker timeout to data-pipeline.queue.timeout. Set queue retry_after (or SQS visibility timeout) above the worker timeout, and configure lease duration and pending-claim TTL above the longest chunk execution. Released jobs count as attempts; the default is 30 attempts with a 5-second release delay.
  • Passwords and other configured raw attributes are masked before quarantine storage. Retry definitions must restore [REDACTED] values from an authoritative source; otherwise those rows remain quarantined. Customize the attribute list for your schema. Database errors stored by built-in sinks contain SQLSTATE, excluding SQL and bindings; custom exception messages and application logs remain the host application's responsibility.
  • Partial custom sinks must return a SinkResult with explicit persisted/failed row indices or add errors to the failed record contexts. Ambiguous counts cause rollback. Bulk sink failures quarantine the entire chunk; ModelSink isolates each row with a savepoint and performs per-row writes.
  • queue() stages the complete input as job payloads before dispatch and therefore has O(N) dispatch memory. Use synchronous streaming or bounded, independently named runs for very large datasets. Resume retains checkpoint metadata in memory. Sources must be replayable with a stable fingerprint and unchanged chunk size. Stop the original workers before resume and use one resume coordinator per run; active chunk leases are rejected.
  • Set a retention policy for completed runs and quarantined data. Deleting a run cascades its checkpoints, errors, and deduplication journal.

Benchmark and stress commands default to pipeline_benchmark_records. Stress tests truncate that dedicated table. Truncating a custom table requires --allow-destructive; benchmarks additionally require --truncate. Generated CSV paths can be configured with --file for benchmarks.

For an existing installation, apply the new additive migration during your normal release:

php artisan migrate --force
php artisan queue:restart

Standalone Package Development & Testing

Run the test suite via Pest:

composer install
composer test
composer analyse
composer validate --strict
composer audit

Real PostgreSQL/Redis crash and concurrency tests use isolated services (database pipeline_test, user pipeline, password pipeline-test):

PIPELINE_INTEGRATION=1 PIPELINE_PG_PORT=5432 PIPELINE_REDIS_PORT=6379 \
  vendor/bin/pest tests/Feature/RedisPostgresConcurrencyTest.php

CI runs the supported PHP/Laravel matrix, static analysis, formatting, dependency audits, and real-service regressions.

Format code via Laravel Pint:

composer format

Security Vulnerabilities

If you discover a security vulnerability within Laravel Data Pipeline, please send an e-mail to Robel at robelsust@gmail.com.

License

Laravel Data Pipeline is open-sourced software licensed under the MIT license.