Weave Code
Code Weaver
Helps Laravel developers discover, compare, and choose open-source packages. See popularity, security, maintainers, and scores at a glance to make better decisions.
Feedback
Share your thoughts, report bugs, or suggest improvements.
Subject
Message

Dataflow Bundle Laravel Package

code-rhapsodie/dataflow-bundle

Symfony bundle for building import/export dataflows: one reader, ordered processing steps (sync/async), and one or more writers. Includes CLI tools and scheduling for jobs, status/result reporting, and support for multiple Doctrine DBAL connections.

View on GitHub
Deep Wiki
Context7

Getting Started

Minimal Setup

  1. Installation:

    composer require code-rhapsodie/dataflow-bundle
    

    Register the bundle in config/bundles.php:

    CodeRhapsodie\DataflowBundle\CodeRhapsodieDataflowBundle::class => ['all' => true],
    
  2. Database Schema: Generate and run migrations for the cr_dataflow_scheduled and cr_dataflow_job tables:

    bin/console code-rhapsodie:dataflow:dump-schema
    
  3. First Dataflow Type: Create a custom DataflowType extending AbstractDataflowType (see example in README). Tag it with coderhapsodie.dataflow.type.

  4. Verify Registration:

    bin/console debug:container --tag coderhapsodie.dataflow.type
    

First Use Case: CSV Import

// src/DataflowType/MyCsvImportType.php
use CodeRhapsodie\DataflowBundle\DataflowType\AbstractDataflowType;
use CodeRhapsodie\DataflowBundle\DataflowType\DataflowBuilder;

class MyCsvImportType extends AbstractDataflowType
{
    protected function buildDataflow(DataflowBuilder $builder, array $options): void
    {
        $builder
            ->setReader(new \CodeRhapsodie\DataflowExemple\Reader\FileReader())
            ->addStep(function ($row) {
                return ['processed' => strtoupper($row[0])];
            })
            ->addWriter(new \CodeRhapsodie\DataflowBundle\DataflowType\Writer\PortWriterAdapter(
                new \Port\FileWriter($options['output'])
            ));
    }

    protected function configureOptions(OptionsResolver $resolver): void
    {
        $resolver->setRequired(['input', 'output']);
    }
}

Run the dataflow:

bin/console code-rhapsodie:dataflow:run MyCsvImportType --input=input.csv --output=output.csv

Implementation Patterns

Core Workflow Patterns

  1. Reader-Step-Writer Pipeline:

    • Reader: Fetch data row-by-row (e.g., CSV, DB query, API).
      $builder->setReader(new \GeneratorFunction(function() {
          yield ['id' => 1, 'name' => 'Item 1'];
          yield ['id' => 2, 'name' => 'Item 2'];
      }));
      
    • Steps: Transform/filter data. Use synchronous (default) or asynchronous (with Generator + scaling):
      // Synchronous step
      $builder->addStep(function ($item) {
          $item['name'] = strtoupper($item['name']);
          return $item;
      });
      
      // Asynchronous step (with 2x scaling)
      $builder->addStep(function ($item): \Generator {
          yield new \Amp\Delayed(100); // Simulate delay
          $item['processed'] = true;
          return $item;
      }, 2);
      
    • Writer: Persist data (e.g., DB, file, API).
      $builder->addWriter(new PortWriterAdapter(new \Port\DatabaseWriter()));
      
  2. Configuration-Driven: Use OptionsResolver to validate/normalize inputs:

    protected function configureOptions(OptionsResolver $resolver): void
    {
        $resolver->setDefaults([
            'batch_size' => 100,
            'timeout' => 30,
        ]);
        $resolver->setRequired(['source', 'destination']);
    }
    
  3. Messenger Integration: Enable concurrent processing via Symfony Messenger:

    # config/packages/code_rhapsodie_dataflow.yaml
    code_rhapsodie_dataflow:
        messenger_mode:
            enabled: true
    

    Configure routing in messenger.yaml:

    routing:
        CodeRhapsodie\DataflowBundle\MessengerMode\JobMessage: async
    

Common Integration Patterns

  1. Doctrine DBAL: Use multiple connections for readers/writers:

    $builder->setReader(new \CodeRhapsodie\DataflowBundle\Reader\DoctrineReader(
        $entityManager->getConnection('secondary')
    ));
    
  2. PortPHP Integration: Adapt PortPHP components (e.g., StreamMergeWriter):

    $builder->addWriter(new PortWriterAdapter(
        new \Port\Writer\StreamMergeWriter()
    ));
    
  3. Scheduled Jobs: Define cron-like schedules via CLI:

    bin/console code-rhapsodie:dataflow:schedule MyDataflowType "0 3 * * *" --param=value
    

    List/enable/disable schedules:

    bin/console code-rhapsodie:dataflow:list-schedules
    bin/console code-rhapsodie:dataflow:enable MyDataflowType
    
  4. Error Handling:

    • Database: Default (store exceptions in cr_dataflow_job).
    • Filesystem: Configure league/flysystem for custom storage:
      code_rhapsodie_dataflow:
          exceptions_mode:
              type: file
              flysystem_service: app.custom_filesystem
      

Advanced Patterns

  1. Dynamic Dataflows: Build dataflows programmatically based on runtime conditions:

    $builder = new DataflowBuilder();
    if ($options['mode'] === 'sync') {
        $builder->addStep(new SyncStep());
    } else {
        $builder->addStep(new AsyncStep(), 4); // Scale async step
    }
    
  2. Step Chaining: Chain steps with dependencies:

    $builder
        ->addStep(function ($item) { /* Step 1 */ return $item; })
        ->addStep(function ($item) use ($dependency) {
            return $dependency->process($item);
        });
    
  3. Logging: Inject the logger into custom components:

    class MyStep implements \CodeRhapsodie\DataflowBundle\DataflowType\StepInterface
    {
        public function __construct(private LoggerInterface $logger) {}
    
        public function __invoke($item)
        {
            $this->logger->info('Processing item', ['item' => $item]);
            return $item;
        }
    }
    

Gotchas and Tips

Common Pitfalls

  1. False Values in Readers:

    • Readers must not yield false (treated as "end of data").
    • Fix: Use null or throw an exception for invalid data.
  2. Async Step Ordering:

    • Steps with scaling factors (>1) may execute out of order.
    • Fix: Set all async steps to scale 1 or use locks.
  3. Messenger Deadlocks:

    • Concurrent jobs may cause race conditions.
    • Fix: Use unique job IDs or implement retry logic.
  4. Database Locks:

    • Long-running readers/writers can block DB connections.
    • Fix: Use transactions or batch processing.
  5. Logger Injection:

    • AbstractDataflowType auto-injects $this->logger, but custom steps must manually inject it.

Debugging Tips

  1. Check Job Status:

    bin/console code-rhapsodie:dataflow:show-last-job MyDataflowType
    
  2. Enable Verbose Logging:

    bin/console code-rhapsodie:dataflow:run MyDataflowType --verbose
    
  3. Inspect Dataflow: Use debug:container to verify tagged services:

    bin/console debug:container CodeRhapsodie\DataflowExemple\DataflowType\MyDataflowType
    
  4. Test Readers/Steps: Isolate components with unit tests:

    public function testReader()
    {
        $reader = new FileReader();
        $data = iterator_to_array($reader->read('test.csv'));
        $this->assertCount(2, $data);
    }
    

Configuration Quirks

  1. Default Connection: Override the DBAL connection in config/packages/code_rhapsodie_dataflow.yaml:

    code_rhapsodie_dataflow:
        dbal_default_connection: my_custom_connection
    
  2. Messenger Transport: Ensure the transport (e.g., async) is configured in messenger.yaml:

    transports:
        async: '%env(MESSENGER_TRANSPORT_DSN)%'
    
  3. Options Resolution:

    • Required options must be set via CLI or config.
    • Defaults are merged after required options.

Extension Points

  1. Custom Writers: Implement WriterInterface for new output formats:

    class ApiWriter implements WriterInterface
    {
        public function prepare() { /* Setup API client */ }
        public function write($item) { /* Send to API */ }
        public function finish() { /* Cleanup */ }
    }
    
  2. **

Weaver

How can I help you explore Laravel packages today?

Conversation history is not saved when not logged in.
Prompt
Add packages to context
No packages found.
terminal42/code-quality-tools
codifyo/ts-generator-bundle
andydefer/laravel-cluster
testo/fiber
mintobit/jobqueue
a4sex/maintenance-bundle
a4sex/entity-date-update
a4sex/client-identifier
a4sex/base-utilites
a4sex/key-value-storage
a4sex/micro-status
chilldev/dependency-injection-extra
datinglibre/datinglibre-app-api
biberltd/corebundle
bricre/symfony-bundle-test
biberltd/logbundle
dominium/http-adapter-bundle
dominium/google-analytics
a4sex/auto-clean-entity
christhompsontldr/laravel-inky