# Streaming / Pipeline / Dataflow Processing

> Every principle in this category is listed as a record.

Page: Ontology · Principles
Canonical: https://banes-lab.com/ontology#architecture-category-streaming-pipeline-dataflow-processing

Listed in [Ontology · Principles](https://banes-lab.com/api/pages/ontology/principles.md), after [Scalability / Performance / Optimization](https://banes-lab.com/ontology/principles/architecture-category-scalability-performance-optimization.md) and before [Plugin / Extensibility / IoC](https://banes-lab.com/ontology/principles/architecture-category-plugin-extensibility-ioc.md).

Every principle in this category is listed as a record. Each record carries its kind, its severity, the scopes it applies at and the layer it lives in, then the edge relations that join it to other records, the records that point back at it, the contracts that answer to it and the tensions it takes part in. The descriptors say how it is violated, detected, measured, repaired and enforced. Where the record carries one, an exemplar shows the shape before and after the principle is applied.

Relations diagram

The relations inside this category.

```mermaid
flowchart LR
n_streaming_architecture["Streaming Architecture"]
n_single_pass_processing["Single-Pass Processing"]
n_pipeline_architecture["Pipeline Architecture"]
n_lazy_evaluation["Lazy Evaluation"]
n_sequential_access["Sequential Access"]
n_forward_only_processing["Forward-Only Processing"]
n_dataflow_architecture["Dataflow Architecture"]
n_stateless_processing["Stateless Processing"]
n_windowing["Windowing"]
n_fan_out_fan_in["Fan-out/Fan-in"]
n_batch_vs_stream["Batch-vs-Stream"]
n_streaming_architecture --> n_single_pass_processing
n_forward_only_processing --> n_single_pass_processing
n_dataflow_architecture --> n_pipeline_architecture
n_windowing --> n_streaming_architecture
```

### Streaming Architecture

- Kind: [style](https://banes-lab.com/records/kind/style.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [contextual](https://banes-lab.com/records/vocabulary/severity-contextual.md)
- Scope: data processing, integration
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
A convention of processing data record by record as it arrives, with backpressure, instead of collecting it first.

Requires
[Event Stream](https://banes-lab.com/records/architecture/event-stream.md), [Backpressure](https://banes-lab.com/records/architecture/backpressure.md)

Reinforces
[Single-Pass Processing](https://banes-lab.com/records/architecture/single-pass-processing.md)

Enables
[Continuous Processing](https://banes-lab.com/records/lexicon/continuous-processing.md)

In tension with
[Ordering/State](https://banes-lab.com/records/lexicon/ordering-state.md), [Batch-Only Processing](https://banes-lab.com/records/lexicon/batch-only-processing.md)

Conflicts with
none

Referenced by
[Event Stream](https://banes-lab.com/records/architecture/event-stream.md), [Windowing](https://banes-lab.com/records/architecture/windowing.md)

Tensions
[Streaming Architecture / Ordering/State](https://banes-lab.com/records/tension/ordering-state-streaming-architecture.md), [Streaming Architecture / Batch-Only Processing](https://banes-lab.com/records/tension/batch-only-processing-streaming-architecture.md)

Violated by
materializing unbounded streams

Detected by
unbounded collection over stream source

Measured by
lag, [throughput](https://banes-lab.com/records/architecture/throughput.md), memory usage

Refactored by
Use Stream Processor, Add Backpressure

Enforced by
load/memory tests

Before

```typescript
const foos = await source.readAll();
const results = foos.map(transformFoo);
await sink.writeAll(results);
```

After

```typescript
for await (const foo of source.stream()) {
await sink.write(transformFoo(foo));
}
```

How it is checked

Checked by
load/memory tests

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Event Stream](https://banes-lab.com/records/architecture/event-stream.md), [Backpressure](https://banes-lab.com/records/architecture/backpressure.md), [Single-Pass Processing](https://banes-lab.com/records/architecture/single-pass-processing.md), [Continuous Processing](https://banes-lab.com/records/lexicon/continuous-processing.md)

Shape it refuses
Not answered

### Single-Pass Processing

- Kind: [principle](https://banes-lab.com/records/kind/principle.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [contextual](https://banes-lab.com/records/vocabulary/severity-contextual.md)
- Scope: algorithm, stream, parser
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
A design rule that large input is read once, with every result it feeds computed in that one pass.

Requires
[Forward-Only State Model](https://banes-lab.com/records/lexicon/forward-only-state-model.md)

Reinforces
[Memory Efficiency](https://banes-lab.com/records/architecture/memory-efficiency.md)

Enables
[Large Input Handling](https://banes-lab.com/records/lexicon/large-input-handling.md)

In tension with
[Global Optimization](https://banes-lab.com/records/lexicon/global-optimization.md), [Multi-Pass Full Materialization](https://banes-lab.com/records/lexicon/multi-pass-full-materialization.md)

Conflicts with
none

Referenced by
[Streaming Architecture](https://banes-lab.com/records/architecture/streaming-architecture.md), [Forward-Only Processing](https://banes-lab.com/records/architecture/forward-only-processing.md)

Tensions
[Single-Pass Processing / Global Optimization](https://banes-lab.com/records/tension/global-optimization-single-pass-processing.md), [Single-Pass Processing / Multi-Pass Full Materialization](https://banes-lab.com/records/tension/multi-pass-full-materialization-single-pass-processing.md)

Violated by
repeated scans over large data where avoidable

Detected by
multiple loops/materializations over same large input

Measured by
pass count, memory use

Refactored by
Fuse Passes, Use Iterator/Accumulator

Enforced by
performance review

Before

```typescript
const names = foos.map(foo => foo.name);
const active = foos.filter(foo => foo.active);
const total = foos.reduce((sum, foo) => sum + foo.count, 0);
```

After

```typescript
const names: string[] = [];
const active: Foo[] = [];
let total = 0;
for (const foo of foos) {
names.push(foo.name);
if (foo.active) active.push(foo);
total += foo.count;
}
```

How it is checked

Checked by
performance review

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Forward-Only State Model](https://banes-lab.com/records/lexicon/forward-only-state-model.md), [Memory Efficiency](https://banes-lab.com/records/architecture/memory-efficiency.md), [Large Input Handling](https://banes-lab.com/records/lexicon/large-input-handling.md)

Shape it refuses
Not answered

### Pipeline Architecture

- Kind: [pattern](https://banes-lab.com/records/kind/pattern.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [recommended](https://banes-lab.com/records/vocabulary/severity-recommended.md)
- Scope: processing, dataflow, build
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
A design pattern that splits processing into ordered stages, each with a declared input and output contract.

Requires
[Stage Contracts](https://banes-lab.com/records/lexicon/stage-contracts.md)

Reinforces
[Composability](https://banes-lab.com/records/architecture/composability.md), [Streaming](https://banes-lab.com/records/lexicon/streaming.md)

Enables
[Stepwise Transformation](https://banes-lab.com/records/lexicon/stepwise-transformation.md)

In tension with
[Error Propagation/Debugging](https://banes-lab.com/records/lexicon/error-propagation-debugging.md)

Conflicts with
[Monolithic Processing Function](https://banes-lab.com/records/lexicon/monolithic-processing-function.md)

Referenced by
[Composability](https://banes-lab.com/records/architecture/composability.md), [Dataflow Architecture](https://banes-lab.com/records/architecture/dataflow-architecture.md)

Tensions
[Pipeline Architecture / Error Propagation/Debugging](https://banes-lab.com/records/tension/error-propagation-debugging-pipeline-architecture.md)

Violated by
one large processor handling all stages

Detected by
long procedural transformation chain

Measured by
stage cohesion, stage contract coverage

Refactored by
Split into Stages, Define Stage Contracts

Enforced by
pipeline tests

Before

```typescript
function processFoo(raw: string) {
const parsed = JSON.parse(raw);
const validated = validateFoo(parsed);
const normalized = normalizeFoo(validated);
return saveFoo(normalized);
}
```

After

```typescript
const fooPipeline = pipeline(
parseJson,
validateWith(FooSchema),
normalizeFoo,
saveFoo,
);
fooPipeline.run(raw);
```

How it is checked

Checked by
pipeline tests

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Stage Contracts](https://banes-lab.com/records/lexicon/stage-contracts.md), [Composability](https://banes-lab.com/records/architecture/composability.md), [Streaming](https://banes-lab.com/records/lexicon/streaming.md), [Stepwise Transformation](https://banes-lab.com/records/lexicon/stepwise-transformation.md)

Shape it refuses
[Monolithic Processing Function](https://banes-lab.com/records/lexicon/monolithic-processing-function.md)

### Lazy Evaluation

- Kind: [approach](https://banes-lab.com/records/kind/approach.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [contextual](https://banes-lab.com/records/vocabulary/severity-contextual.md)
- Scope: computation, collection, stream
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
An approach in which a value is computed only when something consumes it.

Requires
[Deferred Execution Semantics](https://banes-lab.com/records/lexicon/deferred-execution-semantics.md)

Reinforces
[Memory Efficiency](https://banes-lab.com/records/architecture/memory-efficiency.md)

Enables
[Avoiding Unneeded Work](https://banes-lab.com/records/lexicon/avoiding-unneeded-work.md)

In tension with
[Debuggability/Resource Lifetime](https://banes-lab.com/records/lexicon/debuggability-resource-lifetime.md), [Eager Full Materialization](https://banes-lab.com/records/lexicon/eager-full-materialization.md)

Conflicts with
none

Tensions
[Lazy Evaluation / Debuggability/Resource Lifetime](https://banes-lab.com/records/tension/debuggability-resource-lifetime-lazy-evaluation.md), [Lazy Evaluation / Eager Full Materialization](https://banes-lab.com/records/tension/eager-full-materialization-lazy-evaluation.md)

Violated by
computing/materializing unused results

Detected by
eager loading of large unused data

Measured by
avoided work, memory reduction

Refactored by
Use Iterator/Generator, Defer Computation

Enforced by
performance tests

Before

```typescript
const normalized = millionFoos.map(normalizeFoo);
const active = normalized.filter(foo => foo.active);
const firstTen = active.slice(0, 10);
```

After

```typescript
const firstTen = sequence(millionFoos)
.map(normalizeFoo)
.filter(foo => foo.active)
.take(10)
.toArray();
```

How it is checked

Checked by
performance tests

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Deferred Execution Semantics](https://banes-lab.com/records/lexicon/deferred-execution-semantics.md), [Memory Efficiency](https://banes-lab.com/records/architecture/memory-efficiency.md), [Avoiding Unneeded Work](https://banes-lab.com/records/lexicon/avoiding-unneeded-work.md)

Shape it refuses
Not answered

### Sequential Access

- Kind: [pattern](https://banes-lab.com/records/kind/pattern.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [contextual](https://banes-lab.com/records/vocabulary/severity-contextual.md)
- Scope: file, stream, iterator
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
A design pattern that reads a source in its stored order through a cursor or scan, instead of looking records up one by one.

Requires
[Ordered Read Model](https://banes-lab.com/records/lexicon/ordered-read-model.md)

Reinforces
[Memory Efficiency](https://banes-lab.com/records/architecture/memory-efficiency.md)

Enables
[Large Data Processing](https://banes-lab.com/records/lexicon/large-data-processing.md)

In tension with
[Lookup Performance](https://banes-lab.com/records/lexicon/lookup-performance.md), [Random Access Requirement](https://banes-lab.com/records/lexicon/random-access-requirement.md)

Conflicts with
none

Tensions
[Sequential Access / Lookup Performance](https://banes-lab.com/records/tension/lookup-performance-sequential-access.md), [Sequential Access / Random Access Requirement](https://banes-lab.com/records/tension/random-access-requirement-sequential-access.md)

Violated by
random access over stream-only source

Detected by
seek/index assumptions on sequential source

Measured by
access pattern cost

Refactored by
Use Buffer/Index or Stream Sequentially

Enforced by
performance tests

Before

```typescript
for (const id of fooIds) await fooStore.randomRead(id);
```

After

```typescript
for await (const foo of fooStore.scan({ orderBy: "id" })) {
processFoo(foo);
}
```

How it is checked

Checked by
performance tests

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Ordered Read Model](https://banes-lab.com/records/lexicon/ordered-read-model.md), [Memory Efficiency](https://banes-lab.com/records/architecture/memory-efficiency.md), [Large Data Processing](https://banes-lab.com/records/lexicon/large-data-processing.md)

Shape it refuses
Not answered

### Forward-Only Processing

- Kind: [constraint](https://banes-lab.com/records/kind/constraint.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [contextual](https://banes-lab.com/records/vocabulary/severity-contextual.md)
- Scope: stream, parser, iterator
- Aliases: No Backtracking Requirement
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
A rule or precondition that a stream is read once in order, keeping only rolling state, so no decision ever needs to revisit earlier input.

Requires
none

Reinforces
[Single-Pass Processing](https://banes-lab.com/records/architecture/single-pass-processing.md)

Enables
[Streaming Parsers](https://banes-lab.com/records/lexicon/streaming-parsers.md)

In tension with
[Complex Grammar/Global State](https://banes-lab.com/records/lexicon/complex-grammar-global-state.md), [Backtracking Algorithm](https://banes-lab.com/records/lexicon/backtracking-algorithm.md)

Conflicts with
none

Tensions
[Forward-Only Processing / Complex Grammar/Global State](https://banes-lab.com/records/tension/complex-grammar-global-state-forward-only-processing.md), [Forward-Only Processing / Backtracking Algorithm](https://banes-lab.com/records/tension/backtracking-algorithm-forward-only-processing.md)

Violated by
requiring prior/future full data in stream path

Detected by
buffering full stream to look back

Measured by
buffer size, pass count

Refactored by
Add Rolling State, Redesign Parser

Enforced by
memory tests

Before

```typescript
const cursor = fooStream.cursor();
cursor.next();
cursor.previous();
cursor.seek(0);
```

After

```typescript
for await (const foo of fooStream) {
await processFoo(foo);
}
```

How it is checked

Checked by
memory tests

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Single-Pass Processing](https://banes-lab.com/records/architecture/single-pass-processing.md), [Streaming Parsers](https://banes-lab.com/records/lexicon/streaming-parsers.md)

Shape it refuses
Not answered

### Dataflow Architecture

- Kind: [style](https://banes-lab.com/records/kind/style.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [contextual](https://banes-lab.com/records/vocabulary/severity-contextual.md)
- Scope: processing, workflow, stream
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
A convention of arranging processing as a graph of stages connected by explicit data edges, where a stage runs when its inputs arrive.

Requires
[Data Dependencies](https://banes-lab.com/records/lexicon/data-dependencies.md), [Stages](https://banes-lab.com/records/lexicon/stages.md)

Reinforces
[Pipeline Architecture](https://banes-lab.com/records/architecture/pipeline-architecture.md)

Enables
[Parallel/Stream Processing](https://banes-lab.com/records/lexicon/parallel-stream-processing.md)

In tension with
[State Coordination](https://banes-lab.com/records/lexicon/state-coordination.md)

Conflicts with
[Control-Flow-Centric Monolith](https://banes-lab.com/records/lexicon/control-flow-centric-monolith.md)

Tensions
[Dataflow Architecture / State Coordination](https://banes-lab.com/records/tension/dataflow-architecture-state-coordination.md)

Violated by
hidden data dependencies between stages

Detected by
implicit shared state in pipeline

Measured by
data dependency clarity

Refactored by
Make Data Edges Explicit, Split Stages

Enforced by
pipeline contracts

Before

```typescript
controller.runFoo();
controller.runBar();
controller.runBaz();
```

After

```typescript
const graph = dataflow()
.source("foo", fooSource)
.map("bar", "foo", toBar)
.map("baz", "bar", toBaz)
.sink("output", "baz", bazSink);
await graph.run();
```

How it is checked

Checked by
pipeline contracts

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Data Dependencies](https://banes-lab.com/records/lexicon/data-dependencies.md), [Stages](https://banes-lab.com/records/lexicon/stages.md), [Pipeline Architecture](https://banes-lab.com/records/architecture/pipeline-architecture.md), [Parallel/Stream Processing](https://banes-lab.com/records/lexicon/parallel-stream-processing.md)

Shape it refuses
[Control-Flow-Centric Monolith](https://banes-lab.com/records/lexicon/control-flow-centric-monolith.md)

### Stateless Processing

- Kind: [principle](https://banes-lab.com/records/kind/principle.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [recommended](https://banes-lab.com/records/vocabulary/severity-recommended.md)
- Scope: function, stream processor, service
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
A design rule that a processor derives each result only from its input and explicitly passed context.

Requires
[Explicit Inputs](https://banes-lab.com/records/lexicon/explicit-inputs.md), [No Hidden State](https://banes-lab.com/records/lexicon/no-hidden-state.md)

Reinforces
[Scalability](https://banes-lab.com/records/architecture/scalability.md), [Testability](https://banes-lab.com/records/architecture/testability.md)

Enables
[Parallel Processing](https://banes-lab.com/records/lexicon/parallel-processing.md)

In tension with
[Stateful Business Rules](https://banes-lab.com/records/lexicon/stateful-business-rules.md)

Conflicts with
[Stateful Hidden Accumulation](https://banes-lab.com/records/lexicon/stateful-hidden-accumulation.md)

Tensions
[Stateless Processing / Stateful Business Rules](https://banes-lab.com/records/tension/stateful-business-rules-stateless-processing.md)

Violated by
hidden mutable state in processor

Detected by
mutable state across records/requests

Measured by
stateful operator count

Refactored by
Externalize State, Pass State Explicitly

Enforced by
[code review](https://banes-lab.com/records/architecture/code-review.md), [tests](https://banes-lab.com/records/lexicon/tests.md)

Before

```typescript
class FooProcessor {
private previous?: Foo;
process(foo: Foo) {
const result = merge(this.previous, foo);
this.previous = foo;
return result;
}
}
```

After

```typescript
function processFoo(foo: Foo, context: Readonly<FooContext>): FooResult {
return deriveFooResult(foo, context);
}
```

How it is checked

Checked by
code review, tests

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Explicit Inputs](https://banes-lab.com/records/lexicon/explicit-inputs.md), [No Hidden State](https://banes-lab.com/records/lexicon/no-hidden-state.md), [Scalability](https://banes-lab.com/records/architecture/scalability.md), [Testability](https://banes-lab.com/records/architecture/testability.md), [Parallel Processing](https://banes-lab.com/records/lexicon/parallel-processing.md)

Shape it refuses
[Stateful Hidden Accumulation](https://banes-lab.com/records/lexicon/stateful-hidden-accumulation.md)

### Windowing

- Kind: [mechanism](https://banes-lab.com/records/kind/mechanism.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [contextual](https://banes-lab.com/records/vocabulary/severity-contextual.md)
- Scope: data processing, streaming, aggregation
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
A mechanism that groups an unbounded stream into bounded windows of event time and aggregates each window separately.

Requires
[Event Time](https://banes-lab.com/records/lexicon/event-time.md)

Reinforces
[Streaming Architecture](https://banes-lab.com/records/architecture/streaming-architecture.md), [Bounded State](https://banes-lab.com/records/lexicon/bounded-state.md)

Enables
[Bounded Aggregation over Unbounded Streams](https://banes-lab.com/records/lexicon/bounded-aggregation-over-unbounded-streams.md)

In tension with
[Late-Data Handling](https://banes-lab.com/records/lexicon/late-data-handling.md)

Conflicts with
[Unbounded Accumulation](https://banes-lab.com/records/lexicon/unbounded-accumulation.md)

Tensions
[Windowing / Late-Data Handling](https://banes-lab.com/records/tension/late-data-handling-windowing.md)

Violated by
aggregating an unbounded stream into ever-growing state

Detected by
unbounded accumulator over a stream

Measured by
aggregation state growth rate

Refactored by
Aggregate over Windows

Enforced by
streaming design review

Before

```typescript
const total = allFooEvents.reduce((sum, event) => sum + event.value, 0);
```

After

```typescript
for await (const window of fooStream.tumbling({ seconds: 60 })) {
emit(window.start, window.events.reduce((sum, event) => sum + event.value, 0));
}
```

How it is checked

Checked by
streaming design review

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Event Time](https://banes-lab.com/records/lexicon/event-time.md), [Streaming Architecture](https://banes-lab.com/records/architecture/streaming-architecture.md), [Bounded State](https://banes-lab.com/records/lexicon/bounded-state.md), [Bounded Aggregation over Unbounded Streams](https://banes-lab.com/records/lexicon/bounded-aggregation-over-unbounded-streams.md)

Shape it refuses
[Unbounded Accumulation](https://banes-lab.com/records/lexicon/unbounded-accumulation.md)

### Fan-out/Fan-in

- Kind: [pattern](https://banes-lab.com/records/kind/pattern.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [contextual](https://banes-lab.com/records/vocabulary/severity-contextual.md)
- Scope: data processing, parallelism, pipeline
- Aliases: Upward Reduction
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
A design pattern that splits work into independent parts processed in parallel, then merges their results.

Requires
[Independent Work Units](https://banes-lab.com/records/lexicon/independent-work-units.md)

Reinforces
[Parallelism](https://banes-lab.com/records/architecture/parallelism.md), [Throughput](https://banes-lab.com/records/architecture/throughput.md)

Enables
[Parallel Branch Processing](https://banes-lab.com/records/lexicon/parallel-branch-processing.md), [Result Aggregation](https://banes-lab.com/records/lexicon/result-aggregation.md)

In tension with
[Coordination Overhead](https://banes-lab.com/records/lexicon/coordination-overhead.md), [Serial Item Processing](https://banes-lab.com/records/lexicon/serial-item-processing.md)

Conflicts with
none

Tensions
[Fan-out/Fan-in / Coordination Overhead](https://banes-lab.com/records/tension/coordination-overhead-fan-out-fan-in.md), [Fan-out/Fan-in / Serial Item Processing](https://banes-lab.com/records/tension/fan-out-fan-in-serial-item-processing.md)

Violated by
independent items processed strictly one at a time

Detected by
serial loop over parallelizable work

Measured by
parallelism utilization

Refactored by
Fan Out Work, Fan In Results

Enforced by
pipeline design review

Before

```typescript
const report = await buildFullFooReport(foos);
```

After

```typescript
const partials = await fanOut(partition(foos), buildPartialFooReport);
const report = fanIn(partials, mergeFooReports);
```

How it is checked

Checked by
pipeline design review

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Independent Work Units](https://banes-lab.com/records/lexicon/independent-work-units.md), [Parallelism](https://banes-lab.com/records/architecture/parallelism.md), [Throughput](https://banes-lab.com/records/architecture/throughput.md), [Parallel Branch Processing](https://banes-lab.com/records/lexicon/parallel-branch-processing.md), [Result Aggregation](https://banes-lab.com/records/lexicon/result-aggregation.md)

Shape it refuses
Not answered

### Batch-vs-Stream

- Kind: [approach](https://banes-lab.com/records/kind/approach.md)
- Category: [Streaming / Pipeline / Dataflow Processing](https://banes-lab.com/ontology/principles/architecture-category-streaming-pipeline-dataflow-processing.md)
- Severity: [contextual](https://banes-lab.com/records/vocabulary/severity-contextual.md)
- Scope: data processing, latency, architecture
- Layer: [Execution Core](https://banes-lab.com/records/layer/execution-core.md)

Details

Definition
An approach in which the processing model, periodic batch or continuous stream, is chosen by how fresh the data must be.

Requires
[Latency Requirement Clarity](https://banes-lab.com/records/lexicon/latency-requirement-clarity.md)

Reinforces
[Fitness for Purpose](https://banes-lab.com/records/lexicon/fitness-for-purpose.md)

Enables
[Latency-Appropriate Processing Model](https://banes-lab.com/records/lexicon/latency-appropriate-processing-model.md)

In tension with
[Operational Duplication](https://banes-lab.com/records/lexicon/operational-duplication.md)

Conflicts with
[One-Size-Fits-All Processing](https://banes-lab.com/records/lexicon/one-size-fits-all-processing.md)

Tensions
[Batch-vs-Stream / Operational Duplication](https://banes-lab.com/records/tension/batch-vs-stream-operational-duplication.md)

Violated by
low-latency needs served by periodic batch jobs

Detected by
batch cadence mismatched to freshness requirements

Measured by
data-freshness lag vs requirement

Refactored by
Choose Batch or Stream by Latency Need

Enforced by
data architecture review

Before

```typescript
schedule.daily(() => reprocessAllFoos());
```

After

```typescript
fooStream.subscribe(foo => processFoo(foo));
```

How it is checked

Checked by
data architecture review

Population
Every stage, stream source, accumulator and data edge of a pipeline

Freshness
A verdict stands until a stage, its contract or the input volume changes

Refusal
The pipeline test or the memory test fails a stage that materializes unbounded input or hides a data dependency

Observation
Stage contracts read from source, plus lag, throughput and memory measured under load

Evidence
None, because the catalog states this check as a class, so a watched run belongs to each system that adopts it

Authoritative side
The stage contract, which each stage's behavior under load is compared against

Depends on
[Latency Requirement Clarity](https://banes-lab.com/records/lexicon/latency-requirement-clarity.md), [Fitness for Purpose](https://banes-lab.com/records/lexicon/fitness-for-purpose.md), [Latency-Appropriate Processing Model](https://banes-lab.com/records/lexicon/latency-appropriate-processing-model.md)

Shape it refuses
[One-Size-Fits-All Processing](https://banes-lab.com/records/lexicon/one-size-fits-all-processing.md)

## Links to

- [style](https://banes-lab.com/records/kind/style.md)
- [contextual](https://banes-lab.com/records/vocabulary/severity-contextual.md)
- [Execution Core](https://banes-lab.com/records/layer/execution-core.md)
- [Event Stream](https://banes-lab.com/records/architecture/event-stream.md)
- [Backpressure](https://banes-lab.com/records/architecture/backpressure.md)
- [Single-Pass Processing](https://banes-lab.com/records/architecture/single-pass-processing.md)
- [Continuous Processing](https://banes-lab.com/records/lexicon/continuous-processing.md)
- [Ordering/State](https://banes-lab.com/records/lexicon/ordering-state.md)
- [Batch-Only Processing](https://banes-lab.com/records/lexicon/batch-only-processing.md)
- [Windowing](https://banes-lab.com/records/architecture/windowing.md)
- [Streaming Architecture / Ordering/State](https://banes-lab.com/records/tension/ordering-state-streaming-architecture.md)
- [Streaming Architecture / Batch-Only Processing](https://banes-lab.com/records/tension/batch-only-processing-streaming-architecture.md)
- [Throughput](https://banes-lab.com/records/architecture/throughput.md)
- [principle](https://banes-lab.com/records/kind/principle.md)
- [Forward-Only State Model](https://banes-lab.com/records/lexicon/forward-only-state-model.md)
- [Memory Efficiency](https://banes-lab.com/records/architecture/memory-efficiency.md)
- [Large Input Handling](https://banes-lab.com/records/lexicon/large-input-handling.md)
- [Global Optimization](https://banes-lab.com/records/lexicon/global-optimization.md)
- [Multi-Pass Full Materialization](https://banes-lab.com/records/lexicon/multi-pass-full-materialization.md)
- [Streaming Architecture](https://banes-lab.com/records/architecture/streaming-architecture.md)
- [Forward-Only Processing](https://banes-lab.com/records/architecture/forward-only-processing.md)
- [Single-Pass Processing / Global Optimization](https://banes-lab.com/records/tension/global-optimization-single-pass-processing.md)
- [Single-Pass Processing / Multi-Pass Full Materialization](https://banes-lab.com/records/tension/multi-pass-full-materialization-single-pass-processing.md)
- [pattern](https://banes-lab.com/records/kind/pattern.md)
- [recommended](https://banes-lab.com/records/vocabulary/severity-recommended.md)
- [Stage Contracts](https://banes-lab.com/records/lexicon/stage-contracts.md)
- [Composability](https://banes-lab.com/records/architecture/composability.md)
- [Streaming](https://banes-lab.com/records/lexicon/streaming.md)
- [Stepwise Transformation](https://banes-lab.com/records/lexicon/stepwise-transformation.md)
- [Error Propagation/Debugging](https://banes-lab.com/records/lexicon/error-propagation-debugging.md)
- [Monolithic Processing Function](https://banes-lab.com/records/lexicon/monolithic-processing-function.md)
- [Dataflow Architecture](https://banes-lab.com/records/architecture/dataflow-architecture.md)
- [Pipeline Architecture / Error Propagation/Debugging](https://banes-lab.com/records/tension/error-propagation-debugging-pipeline-architecture.md)
- [approach](https://banes-lab.com/records/kind/approach.md)
- [Deferred Execution Semantics](https://banes-lab.com/records/lexicon/deferred-execution-semantics.md)
- [Avoiding Unneeded Work](https://banes-lab.com/records/lexicon/avoiding-unneeded-work.md)
- [Debuggability/Resource Lifetime](https://banes-lab.com/records/lexicon/debuggability-resource-lifetime.md)
- [Eager Full Materialization](https://banes-lab.com/records/lexicon/eager-full-materialization.md)
- [Lazy Evaluation / Debuggability/Resource Lifetime](https://banes-lab.com/records/tension/debuggability-resource-lifetime-lazy-evaluation.md)
- [Lazy Evaluation / Eager Full Materialization](https://banes-lab.com/records/tension/eager-full-materialization-lazy-evaluation.md)
- [Ordered Read Model](https://banes-lab.com/records/lexicon/ordered-read-model.md)
- [Large Data Processing](https://banes-lab.com/records/lexicon/large-data-processing.md)
- [Lookup Performance](https://banes-lab.com/records/lexicon/lookup-performance.md)
- [Random Access Requirement](https://banes-lab.com/records/lexicon/random-access-requirement.md)
- [Sequential Access / Lookup Performance](https://banes-lab.com/records/tension/lookup-performance-sequential-access.md)
- [Sequential Access / Random Access Requirement](https://banes-lab.com/records/tension/random-access-requirement-sequential-access.md)
- [constraint](https://banes-lab.com/records/kind/constraint.md)
- [Streaming Parsers](https://banes-lab.com/records/lexicon/streaming-parsers.md)
- [Complex Grammar/Global State](https://banes-lab.com/records/lexicon/complex-grammar-global-state.md)
- [Backtracking Algorithm](https://banes-lab.com/records/lexicon/backtracking-algorithm.md)
- [Forward-Only Processing / Complex Grammar/Global State](https://banes-lab.com/records/tension/complex-grammar-global-state-forward-only-processing.md)
- [Forward-Only Processing / Backtracking Algorithm](https://banes-lab.com/records/tension/backtracking-algorithm-forward-only-processing.md)
- [Data Dependencies](https://banes-lab.com/records/lexicon/data-dependencies.md)
- [Stages](https://banes-lab.com/records/lexicon/stages.md)
- [Pipeline Architecture](https://banes-lab.com/records/architecture/pipeline-architecture.md)
- [Parallel/Stream Processing](https://banes-lab.com/records/lexicon/parallel-stream-processing.md)
- [State Coordination](https://banes-lab.com/records/lexicon/state-coordination.md)
- [Control-Flow-Centric Monolith](https://banes-lab.com/records/lexicon/control-flow-centric-monolith.md)
- [Dataflow Architecture / State Coordination](https://banes-lab.com/records/tension/dataflow-architecture-state-coordination.md)
- [Explicit Inputs](https://banes-lab.com/records/lexicon/explicit-inputs.md)
- [No Hidden State](https://banes-lab.com/records/lexicon/no-hidden-state.md)
- [Scalability](https://banes-lab.com/records/architecture/scalability.md)
- [Testability](https://banes-lab.com/records/architecture/testability.md)
- [Parallel Processing](https://banes-lab.com/records/lexicon/parallel-processing.md)
- [Stateful Business Rules](https://banes-lab.com/records/lexicon/stateful-business-rules.md)
- [Stateful Hidden Accumulation](https://banes-lab.com/records/lexicon/stateful-hidden-accumulation.md)
- [Stateless Processing / Stateful Business Rules](https://banes-lab.com/records/tension/stateful-business-rules-stateless-processing.md)
- [Code Review](https://banes-lab.com/records/architecture/code-review.md)
- [Tests](https://banes-lab.com/records/lexicon/tests.md)
- [mechanism](https://banes-lab.com/records/kind/mechanism.md)
- [Event Time](https://banes-lab.com/records/lexicon/event-time.md)
- [Bounded State](https://banes-lab.com/records/lexicon/bounded-state.md)
- [Bounded Aggregation over Unbounded Streams](https://banes-lab.com/records/lexicon/bounded-aggregation-over-unbounded-streams.md)
- [Late-Data Handling](https://banes-lab.com/records/lexicon/late-data-handling.md)
- [Unbounded Accumulation](https://banes-lab.com/records/lexicon/unbounded-accumulation.md)
- [Windowing / Late-Data Handling](https://banes-lab.com/records/tension/late-data-handling-windowing.md)
- [Independent Work Units](https://banes-lab.com/records/lexicon/independent-work-units.md)
- [Parallelism](https://banes-lab.com/records/architecture/parallelism.md)
- [Parallel Branch Processing](https://banes-lab.com/records/lexicon/parallel-branch-processing.md)
- [Result Aggregation](https://banes-lab.com/records/lexicon/result-aggregation.md)
- [Coordination Overhead](https://banes-lab.com/records/lexicon/coordination-overhead.md)
- [Serial Item Processing](https://banes-lab.com/records/lexicon/serial-item-processing.md)
- [Fan-out/Fan-in / Coordination Overhead](https://banes-lab.com/records/tension/coordination-overhead-fan-out-fan-in.md)
- [Fan-out/Fan-in / Serial Item Processing](https://banes-lab.com/records/tension/fan-out-fan-in-serial-item-processing.md)
- [Latency Requirement Clarity](https://banes-lab.com/records/lexicon/latency-requirement-clarity.md)
- [Fitness for Purpose](https://banes-lab.com/records/lexicon/fitness-for-purpose.md)
- [Latency-Appropriate Processing Model](https://banes-lab.com/records/lexicon/latency-appropriate-processing-model.md)
- [Operational Duplication](https://banes-lab.com/records/lexicon/operational-duplication.md)
- [One-Size-Fits-All Processing](https://banes-lab.com/records/lexicon/one-size-fits-all-processing.md)
- [Batch-vs-Stream / Operational Duplication](https://banes-lab.com/records/tension/batch-vs-stream-operational-duplication.md)

## Linked from

- [The layer topology](https://banes-lab.com/ontology/schema/the-layer-topology.md)
- [The membership](https://banes-lab.com/ontology/schema/the-membership.md)
