Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,13 @@ The repository is organized into two main Python packages:
- **`sentry_streams_k8s/`** - Kubernetes integration and deployment automation for Sentry Streams (pure Python package)
- See [sentry_streams_k8s/AGENTS.md](./sentry_streams_k8s/AGENTS.md) for package-specific development instructions

## Architecture documents

Overall structure of the streaming platform is in the [README](./README.md).
It contains links to the specific module.

The architecture of the platform itself — subsystems, decisions, invariants, guarantees
and known debt — is in [sentry_streams/docs/architecture](./sentry_streams/docs/architecture/README.md).

### Other Directories

Expand Down
70 changes: 66 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,9 +21,6 @@ The main features are:
- Support for stateful and stateless transformations. The state storage is
provided by the platform rather than being part of the application.

- Distributed execution. The primitives used to build the application can
be distributed on multiple nodes by configuration.

- Hide the Kafka details from the application. Like commit policy and topic
partitioning.

Expand All @@ -34,7 +31,72 @@ The main features are:

- Support for multiple runtimes.

[Streams Documentation](https://getsentry.github.io/streams/)
[Streams User Documentation](https://getsentry.github.io/streams/)

## Architecture of a Streaming Application

The three main components of a streaming application built with this library are:

* The pipeline itself written via the Pipeline DSL. This represents the application logic.

* The runtime. This is the binary that executes the pipeline primitives and processes data.

* The Kubernetes infrastructure to deploy the pipeline (when on Kubernetes).

```mermaid
flowchart TB
Kafka[Kafka]

subgraph StreamingApp["Streaming Application"]
subgraph deployment["Kubernetes deployment"]
subgraph Pipeline["Pipeline"]
direction LR
source[source] --> transform[transform] --> sink[sink]
end

Runtime[Runtime]
Runtime -->|executes| Pipeline

Config[Pipeline ConfigMap]
end

subgraph K8s["Kubernetes Infra"]
Macro["sentry-kube macro"]
Operator[Operator]
end
end
K8s -- manages --> deployment
K8s -- manages --> Config
Kafka --> source
```

Pipeline and runtime architecture docs are [here](./sentry_streams/docs/architecture/README.md)

The pipeline DSL allows the user to define the streaming pipeline through a set
of dataflow primitives chained together. This is a Python DSL now.

The application logic is separate from the infrastructure configuration which
is provided as a separate yaml file. The config file covers aspects like
observability, scale, parallelism, tuning, etc. This is meant to separate
the application logic from infra.

The Runtime consumes the pipeline defined above and the config file, manages
the connectivity with Kafka and executes the processing in a supposedly optimized
way.

The platform is provided as a Python package that contains also a native
portion as the runner. This package (`sentry_streams`) is imported by the application
code. The application defines the pipeline in its own code base and then uses
the runner provided by the library for the execution.

This system provides the infrastructure needed to run the application in Kubernetes.
There are two options:

* A [Sentry Kube](https://github.com/getsentry/sentry-infra-tools#sentry-kube) macro.
This can be used to generate Deployments and Configmap manifests to deploy manually

* An experimental operator that manages a pipeline as a CRD.


## to develop in this repo

Expand Down
4 changes: 4 additions & 0 deletions sentry_streams/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,10 @@ the runner CLI script and the Arroyo adapter to run the application on single no
- **`integration_tests/`** - Integration tests
- **`docs/`** - Sphinx documentation

## Architecture documents

Architecture is described [here](./docs/architecture/README.md)

## Quick Start

**Recommended:** Use the repository root Makefile commands documented in [../AGENTS.md](../AGENTS.md):
Expand Down
2 changes: 1 addition & 1 deletion sentry_streams/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ edition = "2021"
[dependencies]
pyo3 = { version = "0.29.0" }
serde = { version = "1.0", features = ["derive"] }
sentry_arroyo = { version = "2.40.0", features = ["ssl"] }
sentry_arroyo = { version = ">=2.40.0, <2.44.0", features = ["ssl"] }
chrono = "0.4.40"
tracing = "0.1.40"
tracing-subscriber = "0.3.20"
Expand Down
11 changes: 11 additions & 0 deletions sentry_streams/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -6,3 +6,14 @@ is designed to handle real-time unbounded data streams.
This is built primarily to allow the creation of Sentry ingestion pipelines
though the api provided is fully independent from the Sentry product and can
be used to build any streaming application.

## Documentation

- **[Architecture](./docs/architecture/README.md)** — the subsystems, the decisions behind
them, the rules that must hold, and what is unfinished. Start there whether you are building
on the platform or taking ownership of it.
- [Adding a step](./docs/architecture/guides/adding-a-step.md) — the procedure for extending
the DSL and the runtime.
- [Deployment configuration](./docs/reference/deployment-config.md) — the config file format.
- [User documentation](https://getsentry.github.io/streams/) — the published Sphinx docs.
- [AGENTS.md](./AGENTS.md) — development environment, tests, type checking.
189 changes: 189 additions & 0 deletions sentry_streams/docs/architecture/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,189 @@
# Sentry Streams Architecture

This directory describes the architecture of the `sentry_streams` package: the roles of
the main subsystems, the decisions behind them, the rules that must hold, and what we
know is unfinished. The decisions themselves, one per page, are in the
[decision log](./decisions/README.md).

# The model

This is the conceptual model the rest of the documentation assumes. Vocabulary is in
the [glossary](./README.md#glossary).

## The three layers

A streaming application built on this platform is made of three independent layers.
The whole design of the package follows from keeping them separate.

**The pipeline DSL** is how a product team describes what their application does. It is
a pure, declarative description of a dataflow graph: a source, a series of
transformations, one or more sinks. It contains no reference to Kafka brokers, thread
counts, consumer groups, or the engine that will run it. An application is a Python
file that builds a `Pipeline` object and assigns it to a module level variable named
`pipeline`.

**The adapter** turns that description into something executable on a specific runtime.
Each adapter implements a common interface, one method per primitive, and is free to
build whatever runtime-specific structure it wants as those methods are called.

**The runner** is the glue and the entry point. It loads the application file and the
deployment configuration, instantiates the requested adapter, walks the pipeline graph
handing each step to the adapter, and tells the adapter to run.

Alongside these there is **the deployment configuration**: a YAML file carrying
everything the pipeline description deliberately left out. It is described once, in
[Deployment configuration](../reference/deployment-config.md).

```mermaid
flowchart TB
subgraph authored["Authored by the product team"]
app["application.py<br/>(Pipeline DSL)"]
cfg["config.yaml<br/>(deployment config)"]
end

runner["Runner<br/>sentry_streams/runner.py"]

subgraph adapters["Adapters"]
rust["RustArroyoAdapter"]
py["ArroyoAdapter (Python)"]
other["...custom adapter"]
end

runtime["Running consumer"]

app -- "pipeline graph" --> runner
cfg -- "infrastructure config" --> runner
runner -- "translate_step() per step" --> adapters
adapters -- "run()" --> runtime
```

The key property of this split is that the same application file can be deployed
against different brokers, topics, parallelism settings and even different runtimes
without being modified.

A fourth layer, the Kubernetes deployment automation, lives in the `sentry_streams_k8s`
package. The seam between the two is described in
[Deployment configuration](../reference/deployment-config.md#the-seam-with-sentry_streams_k8s).

## The primitives

The DSL is built from a small set of primitives, in a few families:

- **Sources** produce messages into the pipeline. `StreamSource` reads from a Kafka
topic. A pipeline has exactly one root source.
- **Transforms** are 1:1 or 1:N operations: `Map`, `PredicateFilter`, `HeadersFilter`,
`FlatMap`.
- **Reduces** accumulate multiple messages into one over a window: `Batch` for plain
batching, `Aggregate` for windowed, optionally keyed, accumulator based aggregation.
- **Branching steps** split the stream: `Router` sends each message to exactly one
downstream branch, `Broadcast` sends a copy to every branch. Both take fully defined
sub-pipelines built with `branch()`.
- **Sinks** terminate a branch: `StreamSink` produces to Kafka, `GCSSink` writes objects
to GCS, `DevNullSink` discards. Every branch must end in a sink; the runner validates
this before anything is built.
- **Complex steps** are higher level primitives that are syntactic sugar over the simple
ones: `Parser`, `Serializer`, `BatchParser`, `ParquetSerializer`, `Reducer`. See
[Pipeline DSL and runner](./pipeline-dsl-and-runner.md#complex-steps).

Refer to [pipeline.py](../../sentry_streams/pipeline/pipeline.py) to find documentation
on each of the implemented steps.

A pipeline is built by chaining these together. Each step has a name, which is also the
key the deployment config uses to configure that step.

```python
from sentry_kafka_schemas.schema_types.ingest_metrics_v1 import IngestMetric

from sentry_streams.examples.transform_metrics import filter_events, transform_msg
from sentry_streams.pipeline.pipeline import (
Map,
Parser,
PredicateFilter,
Serializer,
StreamSink,
streaming_source,
)

pipeline = streaming_source(name="myinput", stream_name="ingest-metrics")

(
pipeline.apply(Parser[IngestMetric]("parser"))
.apply(PredicateFilter("filter", function=filter_events))
.apply(Map("transform", function=transform_msg))
.apply(Serializer("serializer"))
.sink(StreamSink("mysink", stream_name="transformed-events"))
)
```

Notice what is absent: no broker address, no consumer group, no parallelism, no runtime.
`ingest-metrics` and `transformed-events` are logical stream names the deployment config
can override with the physical topics of a given environment.

This snippet is the canonical example. Other documents refer to it rather than repeating
it.

## Build, then run

Building and running are two distinct phases, and the adapter interface is shaped around
that distinction.

During the **build phase** the runner traverses the pipeline graph from the root source,
calling one adapter method per step through a `RuntimeTranslator`. The traversal follows
data flow rather than a list: each method receives the handle of the stream it attaches
to and returns the handle it produces, and branching steps return one handle per branch.
Nothing executes; the adapter accumulates.

During the **run phase** the runner calls `run()` on the adapter, which blocks until the
consumer shuts down.

That these are two separate functions is an invariant, not an implementation detail —
see [Contracts](./contracts.md#the-runner).

## Multiple adapters

The adapter interface is generic in the type of the stream handle, so adapters can
represent "a stream" however their runtime requires. The runner knows nothing about what
flows through it.

The adapter used in production is the **Rust Arroyo adapter** (`rust_arroyo`), documented
in [The Rust Arroyo adapter](./rust-arroyo/README.md).


# Reading paths

**Evolving the platform**

2. [Adding a step](./guides/adding-a-step.md) — the procedure, with anchors and a skeleton.
3. [Contracts and invariants](./contracts.md) — the rules a new step must not break.
4. The [Rust runtime](./rust-arroyo/README.md) page for the area you are touching.

**Architecture**

1. [Decision log](./decisions/README.md) — the choices this design rests on, and why.
2. [Guarantees and failure modes](./guarantees.md) — what the system promises.


# Glossary

Several of these words are used in more than one sense elsewhere in the industry, and
two of them are used in more than one sense in this repository.

| Term | Meaning here |
| --- | --- |
| **Step** | A node in the pipeline the application author writes: `Map`, `Batch`, `StreamSink`. |
| **Primitive** | A step type the adapter interface has a method for (`source`, `map`, `filter`, `reduce`, `sink`, `router`, `broadcast`, `flat_map`). |
| **Complex step** | A step that is syntactic sugar: it `convert()`s to simple steps before an adapter sees it. `Parser`, `Serializer`, `Reducer`. |
| **Adapter** | The component that turns a pipeline description into something a specific runtime can execute. |
| **Handle** | Whatever an adapter uses to represent "a stream" while the graph is being walked. Opaque to the runner. For the Rust adapter it is a `Route`. |
| **Route** | A source name plus an ordered list of waypoints. Identifies a branch of the pipeline. |
| **Waypoint** | One branch name appended to a route by a `Router` or `Broadcast`. |
| **Operator** (`RuntimeOperator`) | A Rust-side *descriptor* of something to build. One operator may cover several DSL steps. |
| **Strategy** | An Arroyo `ProcessingStrategy`. What an operator is built into at run time. |
| **Chain** / **fusion** | Consecutive `Map` steps accumulated by the adapter and composed into a single function. |
| **Delegate** | A Python object implementing `RustOperatorDelegate`, driven by the Rust `PythonAdapter` strategy. |
| **Segment** (config list) | An entry in `pipeline.segments`, selected with `--segment-id`. The deployment-level unit: one Kubernetes workload per entry. |
| **Segment** (`starts_segment`) | A boundary between chains *within* one process. Marks where fusion stops and a new parallelism setting applies. |
| **Watermark** | A periodic control message carrying accumulated offsets. Commits are driven by watermarks, not by data messages. |
| **Committable** | The set of `(topic, partition) → offset` entries a message or watermark is responsible for. |
| **`RawMessage`** | A payload the runtime understands: bytes, readable from Rust without the GIL. |
| **`PyAnyMessage`** | A payload only the application understands: an opaque reference to a Python object. |
64 changes: 64 additions & 0 deletions sentry_streams/docs/architecture/contracts.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
# Contracts and invariants

Rules that hold across the platform. Each one is load-bearing: something breaks, usually
subtly, if it stops holding. Use this as a review checklist.

Behavioural guarantees to the outside world — delivery semantics, failure policy — are in
[Guarantees](./guarantees.md).

## The runner

| Invariant | Why |
| --- | --- |
| Loading the runtime and running it are separate functions | The Rust CLI embeds a Python interpreter, calls the loading function to obtain the built runtime, and calls `run()` itself from Rust. Merging them breaks that entry point |
| Exceptions from product code propagate | A Sentry SDK initialised by the application has to be able to capture them |
| Every branch terminates in a sink, validated before anything is built | A dangling branch is nearly always a bug; catching it after startup is far more expensive |

## Configuration

| Invariant | Why |
| --- | --- |
| `override_config()` then `validate()`, called by the adapter at translation time, and nowhere else | Before it, the step holds what the application wrote; after it, what this deployment wants. A step valid as written may be invalid once overridden, so validation must follow the override |
| No other component reads config keys off a step | One place to look when a value is wrong |
| Configuration is plain data, re-applied per process | It has to survive being pickled into a multiprocessing worker. A live backend object would not |
| Only non-semantic values are overridable | Functions, graph shape and types belong to the application; topics, brokers, consumer groups, batch sizes and parallelism belong to the deployment |

## The adapter

| Invariant | Why |
| --- | --- |
| The fused function of a closed chain is checked for picklability, even when multiprocessing is off | Otherwise a pipeline works in a single-process deployment and fails when someone enables a pool in production |
| Parallelism is configured only on a step that starts a segment | Parallelism belongs to the segment; a step in the middle of a chain is not its own execution unit |
| Nothing executes during the build phase | The graph traversal only accumulates; the runtime is assembled in `run()` |
| Operator descriptors are kept, not consumed | The strategy chain is rebuilt from them on every rebalance |

## The message model

| Invariant | Why |
| --- | --- |
| A strategy checks the route first and forwards untouched on mismatch, without taking the GIL | Every message physically traverses every strategy; this is what makes that affordable |
| A strategy forwards watermarks unless it buffers messages, in which case it holds them until the buffered work is out | A committed offset must be a processed offset |
| A `PyWatermark` is never submitted back into a Python operator | Directional conversion invariant of `WatermarkMessage` |
| A `PyWatermark` is never handed to the Kafka sink | Same |
| A step handles every payload variant: supports what it can, panics loudly on what it cannot | Silently forwarding an unsupported payload corrupts data downstream instead of failing at the bug |
| Payloads are not converted unless they are read | An opaque payload is a pointer move; a converted one is a copy plus a GIL acquisition |
| Messages are immutable; replacing a payload produces a new message | Enforced by Rust, a convention in Python |

## Extending the DSL

Adding a primitive is a **breaking change to every adapter**, by construction: the adapter
interface has one abstract method per primitive, and `RuntimeTranslator` dispatches on a
closed `StepType` enum. That is a deliberate trade.

What it obliges you to do:

1. Every adapter must implement the new method. An adapter that will not support the
primitive raises `NotImplementedError` — that is the sanctioned way to decline, and is
how `flat_map` is handled today in the Rust adapter.
2. `RuntimeTranslator.translate_step` must gain a branch, or `assert_never` fails the type
check.
3. If the primitive is not a map, its adapter methods must close the chain (see above).
4. Prefer a `ComplexStep` if the behaviour can be composed from existing primitives. It
costs no adapter changes and adapters may still override it natively.

The procedure is in [Adding a step](./guides/adding-a-step.md).
Loading
Loading