# Flow Producer: Parent and Child Jobs in Bun

Build durable job graphs in bunqueue: children run first, then the parent reads their results. Flows are atomic in SQLite and PostgreSQL modes.

Canonical: https://bunqueue.dev/guide/flow/

---

import { Tabs, TabItem } from '@astrojs/starlight/components';

<div class="bq-wrap bq-hero">
  <span class="bq-eyebrow">guide · flow producer</span>
  <h1 class="bq-hero-h1 bq-bench-h1">Jobs that <em>wait for each other.</em></h1>
  <p class="bq-hero-sub">Some work only makes sense in order: resize every image, then build the album; charge each line, then close the invoice. A flow declares that shape once and bunqueue holds the parent until its children are done.</p>
</div>

A flow is a tree of jobs. Children are queued immediately, the parent stays blocked until every child has completed, and then runs with access to what they returned.

## Quick Start

The most common shape: a parent job that waits for its children. Children run first, then the parent runs with access to their results:

<Tabs syncKey="lang">
<TabItem label="Bun">

```typescript
import { FlowProducer, Worker } from 'bunqueue/client';

type ReportData = {
  month?: string;
  source?: 'sales' | 'costs';
};

const flow = new FlowProducer({ embedded: true });
const rowsBySource = { sales: [120, 80], costs: [50, 20] };

const worker = new Worker<ReportData>(
  'reports',
  async (job) => {
    if (job.name === 'build-report') {
      const values = await job.getChildrenValues<{ rows: number[] }>();
      return { report: Object.values(values) };
    }
    if (!job.data.source) throw new Error('source is required');
    return { rows: rowsBySource[job.data.source] };
  },
  { embedded: true }
);

const node = await flow.add<ReportData>({
  name: 'build-report',
  queueName: 'reports',
  data: { month: '2026-01' },
  children: [
    { name: 'fetch-sales', queueName: 'reports', data: { source: 'sales' } },
    { name: 'fetch-costs', queueName: 'reports', data: { source: 'costs' } },
  ],
});

const result = await node.job.waitUntilFinished(null, 10_000);
console.log(result);

await worker.close();
await flow.close();
```

</TabItem>
<TabItem label="Node.js / Deno">

```typescript
import { FlowProducer, Worker } from 'bunqueue-client';

const flow = new FlowProducer();

await flow.add({
  name: 'build-report',
  queueName: 'reports',
  data: { month: '2026-01' },
  children: [
    { name: 'fetch-sales', queueName: 'reports', data: { source: 'sales' } },
    { name: 'fetch-costs', queueName: 'reports', data: { source: 'costs' } },
  ],
});

new Worker('reports', async (job) => {
  if (job.name === 'build-report') {
    // Children have completed; read their results
    const values = await job.getChildrenValues();
    return { report: Object.values(values) };
  }
  return { rows: await fetchData(job.data.source) };
});
```

</TabItem>
<TabItem label="Python">

```python
from bunqueue import FlowProducer, Worker

flow = FlowProducer()

flow.add({
    "name": "build-report",
    "queueName": "reports",
    "data": {"month": "2026-01"},
    "children": [
        {"name": "fetch-sales", "queueName": "reports", "data": {"source": "sales"}},
        {"name": "fetch-costs", "queueName": "reports", "data": {"source": "costs"}},
    ],
})

def process(job):
    if job.name == "build-report":
        # Children have completed; read their results
        values = job.get_children_values()
        return {"report": list(values.values())}
    return {"rows": fetch_data(job.data["source"])}

Worker("reports", process).run()
```

</TabItem>
<TabItem label="PHP">

```php
use Bunqueue\FlowProducer;
use Bunqueue\Queue;
use Bunqueue\Worker;

$flow = new FlowProducer();

$flow->add([
    'name' => 'build-report', 'queueName' => 'reports',
    'data' => ['month' => '2026-01'],
    'children' => [
        ['name' => 'fetch-sales', 'queueName' => 'reports', 'data' => ['source' => 'sales']],
        ['name' => 'fetch-costs', 'queueName' => 'reports', 'data' => ['source' => 'costs']],
    ],
]);

$queue = new Queue('reports');
$worker = new Worker('reports', function (Bunqueue\Job $job) use ($queue) {
    if ($job->name() === 'build-report') {
        // Children have completed; read their results
        $values = $queue->getChildrenValues($job->id());
        return ['report' => array_values($values)];
    }
    return ['rows' => fetchData($job->data()['source'])];
});
$worker->run();
```

</TabItem>
<TabItem label="Go">

```go
flow := bunqueue.NewFlowProducer(bunqueue.Options{})
defer flow.Close()

node, err := flow.Add(bunqueue.FlowJob{
    Name: "build-report", QueueName: "reports",
    Data: map[string]any{"month": "2026-01"},
    Children: []bunqueue.FlowJob{
        {Name: "fetch-sales", QueueName: "reports", Data: map[string]any{"source": "sales"}},
        {Name: "fetch-costs", QueueName: "reports", Data: map[string]any{"source": "costs"}},
    },
})

queue := bunqueue.NewQueue("reports", bunqueue.Options{})
worker := bunqueue.NewWorker("reports", func(job *bunqueue.Job) (any, error) {
    if job.Name() == "build-report" {
        // Children have completed; read their results
        values, err := queue.GetChildrenValues(job.ID())
        if err != nil {
            return nil, err
        }
        return map[string]any{"report": values}, nil
    }
    return fetchData(job.Data()["source"].(string))
}, bunqueue.WorkerOptions{})
worker.Run()
```

</TabItem>
<TabItem label="Rust">

```rust
use bunqueue_client::{ConnectionOptions, FlowJob, FlowProducer, JobOptions, Value};

let flow = FlowProducer::new(ConnectionOptions::default());
let child = |name: &str, source: &str| FlowJob {
    name: name.into(),
    queue_name: "reports".into(),
    data: Value::Map(vec![(Value::from("source"), Value::from(source))]),
    options: JobOptions::default(),
    children: vec![],
};
let node = flow.add(FlowJob {
    name: "build-report".into(),
    queue_name: "reports".into(),
    data: Value::Map(vec![(Value::from("month"), Value::from("2026-01"))]),
    options: JobOptions::default(),
    children: vec![child("fetch-sales", "sales"), child("fetch-costs", "costs")],
})?;
```

_Reading children results (`GetChildrenValues`) has no typed helper in the Rust SDK yet; use the documented wire protocol until it is added._

</TabItem>
<TabItem label="Elixir">

```elixir
flow = Bunqueue.FlowProducer.new()

{:ok, node} =
  Bunqueue.FlowProducer.add(flow, %{
    name: "build-report",
    queue: "reports",
    data: %{month: "2026-01"},
    children: [
      %{name: "fetch-sales", queue: "reports", data: %{source: "sales"}},
      %{name: "fetch-costs", queue: "reports", data: %{source: "costs"}}
    ]
  })
```

_Reading children results (`GetChildrenValues`) has no typed helper in the Elixir SDK yet; use the documented wire protocol until it is added._

</TabItem>
</Tabs>

Flow creation is one broker-side transaction in the Bun package and all six
current external SDKs: `addBulk` commits every tree or none, and workers cannot
see a leaf before the full graph exists. Previously published SDK versions that
compose `PUSH` and `UpdateParent` remain compatible with the server, but their
already-sent requests cannot gain `PUSHF` all-or-nothing visibility.

:::tip
In the Bun package, FlowProducer works in TCP mode too: pass `connection: { port: 6789 }` instead of `embedded: true`.
:::

### Creation guarantees and limits

The Bun producer validates the complete graph before sending it and the broker
validates it again. With a SQLite `dataPath`, one immediate transaction commits
the graph before it is published to workers, even when individual nodes omit
`durable: true`. PostgreSQL likewise admits the complete graph, dependency
edges, queue registry changes, and durable events in one database transaction
before any broker can claim a leaf. Without either persistent backend,
embedded mode is intentionally memory-only: creation is still atomically
visible to workers, but a process crash cannot recover it.

- A flow may contain at most 10,000 jobs, at most 10 MB of data per job and
  64 MB across the batch. A root is depth 0; descendants may be at most 100
  edges below it.
- `jobId` is allowed, but cannot be empty or contain `:`. Reusing any existing
  or retained flow ID—including durable job/DLQ rows, completion or timeout
  tombstones, retained results, and IDs still referenced by a waiting
  parent—rejects the whole request in either persistent backend.
- `name` inside user data is preserved independently from the job's own name.
  Keys beginning with `__` are reserved for engine-owned flow metadata.
- `repeat`, `deduplication`, `debounce`, and `opts.parent` are rejected inside
  an atomic flow because their independent lifetime/ownership semantics cannot
  participate safely in the graph transaction.
- `opts.group` is preserved for every node, including `id`, BullMQ Pro
  intra-group `priority`, and `maxSize`. Group options are validated with the
  rest of the graph; an invalid option or full group rejects every node before
  either SQLite or PostgreSQL publishes a partial flow.
- The four child-failure policies are mutually exclusive.

Wire values are checked at runtime too: IDs must be strings, link fields must be
string arrays, booleans must actually be booleans, and parent metadata must
match both sides of every edge. These checks happen before any queue counter,
heap, dependency index, or selected-backend row changes.

:::note
When a child uses `removeOnComplete`, bunqueue retains a payload-free completion
proof so a restart cannot strand its parent. Unreferenced proofs follow
`maxCompletedJobs`; a proof referenced by a waiting parent stays pinned until
that parent's durable promotion or removal releases the final dependency edge.
It does not retain the child Job or result. Once a parent is promoted, its
persisted ready state remains authoritative even after the proof expires. Use
retained children—not `removeOnComplete`—when a parent must read results after
a broker restart.
:::

## Where to go next

|                                                   |                                                       |
| ------------------------------------------------- | ----------------------------------------------------- |
| [Flow Patterns](/guide/flow/patterns/)            | Chains, fan-in, trees, reading child results, options |
| [Flow Failure Handling](/guide/flow/failures/)    | What a parent does when a child dies for good         |
| [Flow Producer Reference](/guide/flow/reference/) | Every producer method, job helper and step field      |