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.
Installation:
composer require code-rhapsodie/dataflow-bundle
Register the bundle in config/bundles.php:
CodeRhapsodie\DataflowBundle\CodeRhapsodieDataflowBundle::class => ['all' => true],
Database Schema:
Generate and run migrations for the cr_dataflow_scheduled and cr_dataflow_job tables:
bin/console code-rhapsodie:dataflow:dump-schema
First Dataflow Type:
Create a custom DataflowType extending AbstractDataflowType (see example in README). Tag it with coderhapsodie.dataflow.type.
Verify Registration:
bin/console debug:container --tag coderhapsodie.dataflow.type
// 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
Reader-Step-Writer Pipeline:
$builder->setReader(new \GeneratorFunction(function() {
yield ['id' => 1, 'name' => 'Item 1'];
yield ['id' => 2, 'name' => 'Item 2'];
}));
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);
$builder->addWriter(new PortWriterAdapter(new \Port\DatabaseWriter()));
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']);
}
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
Doctrine DBAL: Use multiple connections for readers/writers:
$builder->setReader(new \CodeRhapsodie\DataflowBundle\Reader\DoctrineReader(
$entityManager->getConnection('secondary')
));
PortPHP Integration:
Adapt PortPHP components (e.g., StreamMergeWriter):
$builder->addWriter(new PortWriterAdapter(
new \Port\Writer\StreamMergeWriter()
));
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
Error Handling:
cr_dataflow_job).league/flysystem for custom storage:
code_rhapsodie_dataflow:
exceptions_mode:
type: file
flysystem_service: app.custom_filesystem
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
}
Step Chaining: Chain steps with dependencies:
$builder
->addStep(function ($item) { /* Step 1 */ return $item; })
->addStep(function ($item) use ($dependency) {
return $dependency->process($item);
});
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;
}
}
False Values in Readers:
false (treated as "end of data").null or throw an exception for invalid data.Async Step Ordering:
>1) may execute out of order.1 or use locks.Messenger Deadlocks:
Database Locks:
Logger Injection:
AbstractDataflowType auto-injects $this->logger, but custom steps must manually inject it.Check Job Status:
bin/console code-rhapsodie:dataflow:show-last-job MyDataflowType
Enable Verbose Logging:
bin/console code-rhapsodie:dataflow:run MyDataflowType --verbose
Inspect Dataflow:
Use debug:container to verify tagged services:
bin/console debug:container CodeRhapsodie\DataflowExemple\DataflowType\MyDataflowType
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);
}
Default Connection:
Override the DBAL connection in config/packages/code_rhapsodie_dataflow.yaml:
code_rhapsodie_dataflow:
dbal_default_connection: my_custom_connection
Messenger Transport:
Ensure the transport (e.g., async) is configured in messenger.yaml:
transports:
async: '%env(MESSENGER_TRANSPORT_DSN)%'
Options Resolution:
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 */ }
}
**
How can I help you explore Laravel packages today?