Skip to content
Open
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
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -170,3 +170,4 @@ cython_debug/

# Generated version files
*/_version.py
docs/message-data-reader-writer-proposal.md
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,7 @@ Optional integrations for different cloud providers can be installed using `plug

Support for parallelisation and hyperparameter optimisation can be installed using `plugboard[ray]`.

Additional optional extras: `plugboard[llm]` for LLM components, `plugboard[redis]` for Redis-based connectors, `plugboard[omq]` for the pyomq backend for ZMQ connectors, and `plugboard[websockets]` for WebSocket I/O.
Additional optional extras: `plugboard[llm]` for LLM components, `plugboard[redis]` for Redis-based connectors, `plugboard[omq]` for the pyomq backend for ZMQ connectors, `plugboard[websockets]` for WebSocket I/O, and `plugboard[gcp-pubsub]`, `plugboard[aws-messaging]` or `plugboard[kafka]` for [message data components](https://docs.plugboard.dev/usage/message-data/).

## ⚡ Quickstart with AI

Expand Down
20 changes: 20 additions & 0 deletions docs/usage/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,26 @@ Plugboard can make use of a message broker for data exchange between components
| `RABBITMQ_URL` | URL for RabbitMQ AMQP message broker (must include credentials if required) | |
| `REDIS_URL` | URL for Redis message broker (must include credentials if required) | |

These brokers carry data *between components* within a run. To read or write records
through an external broker as part of the model itself, see
[Message Data](message-data.md).

### Message data brokers

The message data components ([`MessageDataReader`][plugboard.library.MessageDataReader] /
[`MessageDataWriter`][plugboard.library.MessageDataWriter] and their broker
implementations) take connection details as constructor arguments, falling back to the
evironment below when they are not supplied:

| Option Name | Description | Default Value |
|---------------------------|----------------------------------------|---------------|
| `GCP_PUBSUB_PROJECT_ID` | GCP project for PubSub topics and subscriptions | |
| `AWS_REGION` | AWS region for SQS and SNS clients | |
| `KAFKA_BOOTSTRAP_SERVERS` | Kafka broker address(es) | |

Note that `AWS_REGION` is also read by the AWS SDK itself, so setting it affects both
Plugboard's defaults and the credential/endpoint resolution of the underlying client.

## Job ID

Each plugboard run has a unique job ID associated with it. This is used to: track state for each run; and separate data messages between runs when using a message broker. Typically, a run would be started without explicitly setting the job ID, in which case a unique job ID will be created automatically. However, there are instances when it may be desirable to specify the job ID, such as stopping a run and resuming the same run later with the existing persisted state. In these scenarios the job ID can be set with the below environment variable which will then be used by any `StateBackend`, `Process` and `Component` while the value is set.
Expand Down
150 changes: 150 additions & 0 deletions docs/usage/message-data.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
# Message Data

Message data components read and write **records** through a pub/sub message broker,
rather than through a channel wired between components. They are the broker equivalent
of [`FileReader`][plugboard.library.FileReader] / [`FileWriter`][plugboard.library.FileWriter]:
each message carries one record, and the fields named by `field_names` become the
component's outputs (for a reader) or inputs (for a writer).

| Class | Broker | Reads from | Writes to |
|---|---|---|---|
| [`GCPPubSubDataReader`][plugboard.library.GCPPubSubDataReader] | Google Cloud PubSub | a subscription | |
| [`GCPPubSubDataWriter`][plugboard.library.GCPPubSubDataWriter] | Google Cloud PubSub | | a topic |
| [`AWSSQSDataReader`][plugboard.library.AWSSQSDataReader] | AWS SQS | a queue | |
| [`AWSSNSDataWriter`][plugboard.library.AWSSNSDataWriter] | AWS SNS | | a topic |
| [`KafkaDataReader`][plugboard.library.KafkaDataReader] | Apache Kafka | a topic + consumer group | |
| [`KafkaDataWriter`][plugboard.library.KafkaDataWriter] | Apache Kafka | | a topic |

All six build on the [`MessageDataReader`][plugboard.library.MessageDataReader] and
[`MessageDataWriter`][plugboard.library.MessageDataWriter] base classes, which handle
connecting, reconnection with backoff, batching and acknowledgment.

## Installing

Each broker uses its own client library, supplied as an extra:

```shell
pip install "plugboard[gcp-pubsub]" # google-cloud-pubsub
pip install "plugboard[aws-messaging]" # aiobotocore
pip install "plugboard[kafka]" # aiokafka
```

Constructing a component without its extra raises an `ImportError` naming the extra to
install.

## Reading and writing records

Messages are expected to carry a JSON object per record, whose keys match
`field_names`:

```python
import asyncio

from plugboard.component import Component, IOController as IO
from plugboard.library import AWSSNSDataWriter, GCPPubSubDataReader
from plugboard.process import LocalProcess
from plugboard.schemas import ConnectorSpec


class Bucket(Component):
"""Turns a raw measurement into a bucketed label before it leaves the model."""

io = IO(inputs=["x", "y"], outputs=["bucket"])

async def step(self) -> None:
self.bucket = "high" if self.y > 50 else "low"


async def main() -> None:
reader = GCPPubSubDataReader(
name="reader",
subscription_id="measurements-pull",
project_id="plugboard-dev", # Optional; defaults to GCP_PUBSUB_PROJECT_ID
field_names=["x", "y"],
chunk_size=20,
)
bucket = Bucket(name="bucket")
writer = AWSSNSDataWriter(
name="writer",
topic_arn="arn:aws:sns:eu-west-1:123456789012:measurements",
region="eu-west-1", # Optional; defaults to AWS_REGION
field_names=["bucket"],
chunk_size=20,
)

process = LocalProcess(
name="measurements",
components=[reader, bucket, writer],
connectors=[
ConnectorSpec(source="reader.x", target="bucket.x"),
ConnectorSpec(source="reader.y", target="bucket.y"),
ConnectorSpec(source="bucket.bucket", target="writer.bucket"),
],
)
await process.run()


asyncio.run(main())
```

The components are ordinary `Component`s, so they also work in a `RayProcess` and with
any connector.

## Behaviour worth knowing

**A reader waits; it does not finish.** An empty poll from a broker means "nothing yet",
not "no more data" — unlike a file, which ends. A reader started before its producer
keeps running and picks up messages as they arrive. To end a reader's stream, the broker
must report the source as gone (for example a deleted PubSub subscription or SQS queue),
which raises `IOStreamClosedError` like other data readers do. Stop the process to end a
run that would otherwise wait forever.

**Messages are acknowledged after they are processed.** A batch fetched from the broker
is acknowledged only once every record in it has been published downstream, so a crash
mid-batch leaves the unread messages to be redelivered (at-least-once). For SQS,
"acknowledge" means deleting the messages; for Kafka it means committing the offsets of
the processed records, per partition.

**Batching is bounded by the broker.** `chunk_size` sets how many messages are fetched
per poll (and per send for writers). SQS delivers at most 10 messages per call, so a
larger `chunk_size` is capped there and a warning is logged.

**Failures are retried with backoff.** Transient broker errors trigger a reconnect and a
retry, doubling the delay up to a cap. Permanent errors — such as `AccessDenied` or a
deleted topic — fail immediately instead of burning the retries. Both the classification
and the policy are handled by the base classes:

```python
from plugboard.library import KafkaDataReader
from plugboard.utils.retry import RetryPolicy

reader = KafkaDataReader(
name="reader",
topic="measurements",
group_id="plugboard",
field_names=["x", "y"],
bootstrap_servers="localhost:9092", # Optional; defaults to KAFKA_BOOTSTRAP_SERVERS
retry_policy=RetryPolicy(max_retries=5, base_delay=0.5, max_delay=30.0),
idle_poll_delay=0.5, # Pause between empty polls, if the broker returns immediately
)
```

**Message encoding.** By default each record is JSON-encoded (`parse_json=True`). With
`parse_json=False` a writer sends only the first field's value, and a reader exposes the
raw payload under a `data` field — useful for single-value or binary messages.

## Configuration

Connection details can be passed explicitly or read from the environment, which lets the
same model run against a different account or region without code changes. Explicit
arguments always win.

| Option Name | Description | Used by |
|---------------------------|------------------------------------------|---------|
| `GCP_PUBSUB_PROJECT_ID` | Default GCP project for PubSub topics | `GCPPubSubDataReader`, `GCPPubSubDataWriter` |
| `AWS_REGION` | Default AWS region for SQS/SNS clients | `AWSSQSDataReader`, `AWSSNSDataWriter` |
| `KAFKA_BOOTSTRAP_SERVERS` | Default Kafka broker address(es) | `KafkaDataReader`, `KafkaDataWriter` |

Credentials themselves are never configured through Plugboard: the clients use their
normal provider mechanisms (Application Default Credentials for PubSub, the standard AWS
credential chain, SASL settings for Kafka).
1 change: 1 addition & 0 deletions mkdocs.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,7 @@ nav:
- Event-driven models: examples/tutorials/event-driven-models.md
- Tuning a process: examples/tutorials/tuning-a-process.md
- Configuration: usage/configuration.md
- Message Data: usage/message-data.md
- AI-Assisted Development: usage/ai.md
- Topics: usage/topics.md
- Demos:
Expand Down
24 changes: 24 additions & 0 deletions plugboard/exceptions/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -116,3 +116,27 @@ class ProcessStatusError(Exception):
"""Raised when a `Process` is in an invalid state for the requested operation."""

pass


class MessageBrokerError(Exception):
"""Base exception for message broker errors."""

pass


class MessageBrokerConnectionError(MessageBrokerError):
"""Raised when connection to a message broker fails."""

pass


class MessageBrokerTransientError(MessageBrokerError):
"""Raised on transient message broker errors (eligible for retry)."""

pass


class MessageBrokerPermanentError(MessageBrokerError):
"""Raised on permanent message broker errors (not eligible for retry)."""

pass
17 changes: 15 additions & 2 deletions plugboard/library/__init__.py
Original file line number Diff line number Diff line change
@@ -1,20 +1,33 @@
"""Provides implementations of Plugboard objects for use in user models."""

from .aws_messaging_io import AWSSNSDataWriter, AWSSQSDataReader
from .data_reader import DataReader
from .data_writer import DataWriter
from .file_io import FileReader, FileWriter
from .gcp_pubsub_io import GCPPubSubDataReader, GCPPubSubDataWriter
from .kafka_io import KafkaDataReader, KafkaDataWriter
from .llm import LLMChat, LLMImageProcessor
from .message_reader import MessageDataReader
from .message_writer import MessageDataWriter
from .sql_io import SQLReader, SQLWriter
from .websocket_io import WebsocketBase, WebsocketReader, WebsocketWriter


__all__ = [
"AWSSQSDataReader",
"AWSSNSDataWriter",
"DataReader",
"DataWriter",
"LLMChat",
"LLMImageProcessor",
"FileReader",
"FileWriter",
"GCPPubSubDataReader",
"GCPPubSubDataWriter",
"KafkaDataReader",
"KafkaDataWriter",
"LLMChat",
"LLMImageProcessor",
"MessageDataReader",
"MessageDataWriter",
"SQLReader",
"SQLWriter",
"WebsocketBase",
Expand Down
Loading
Loading