Metadata-Version: 2.4
Name: fastmesh-kafka
Version: 0.1.0
Summary: Kafka adapter for FastMesh — KafkaSource and KafkaSink via aiokafka
License: MIT
Requires-Python: >=3.11
Requires-Dist: aiokafka>=0.11
Requires-Dist: fastmesh>=0.2
Provides-Extra: dev
Requires-Dist: fastmesh>=0.2; extra == 'dev'
Requires-Dist: pytest-asyncio>=0.24; extra == 'dev'
Requires-Dist: pytest-cov>=5.0; extra == 'dev'
Requires-Dist: pytest>=8.0; extra == 'dev'
Requires-Dist: ruff>=0.4; extra == 'dev'
Description-Content-Type: text/markdown

# fastmesh-kafka

Kafka adapter for [FastMesh](https://github.com/your-org/fastmesh) —
`KafkaSource` and `KafkaSink` via [aiokafka](https://aiokafka.readthedocs.io/).

## Install

```bash
pip install fastmesh-kafka
```

## Usage

```python
from fastmesh import Bus
from fastmesh_kafka import KafkaSource, KafkaSink

source = KafkaSource(
    topics=["orders"],
    bootstrap_servers="localhost:9092",
    group_id="pipeline-consumer",
    auto_offset_reset="earliest",
)

sink = KafkaSink(
    topic="processed-orders",
    bootstrap_servers="localhost:9092",
)

await Bus().source(source).sink(sink).run()
```

## KafkaSource

| Parameter | Default | Description |
|-----------|---------|-------------|
| `topics` | *(required)* | List of topic names |
| `bootstrap_servers` | `"localhost:9092"` | Broker addresses |
| `group_id` | `"fastmesh"` | Consumer group ID |
| `auto_offset_reset` | `"latest"` | `"earliest"` or `"latest"` |
| `enable_auto_commit` | `True` | Auto-commit offsets |
| `value_deserializer` | JSON decode | `(bytes) -> dict` |

Each Kafka record → one `Message`. Kafka metadata in `msg.metadata`:
`kafka_topic`, `kafka_partition`, `kafka_offset`, `kafka_key` (if present).

## KafkaSink

| Parameter | Default | Description |
|-----------|---------|-------------|
| `topic` | *(required)* | Target topic name |
| `bootstrap_servers` | `"localhost:9092"` | Broker addresses |
| `key_fn` | `msg.id.encode()` | `(Message) -> bytes \| None` |
| `value_serializer` | JSON encode | `(dict) -> bytes` |

The producer is lazily initialized on first `send()` and reused. Call
`await sink.stop()` on shutdown to flush and close.

## With FastMesh plugins

```python
from fastmesh.plugins.retry import Retry
from fastmesh.plugins.dlq import DeadLetterQueue, MemoryDLQStore

store = MemoryDLQStore()
bus = Bus().source(
    KafkaSource(topics=["orders"], bootstrap_servers="kafka:9092")
).sink(
    DeadLetterQueue(
        Retry(KafkaSink("processed", bootstrap_servers="kafka:9092"), max_retries=3),
        store=store,
    )
)
```

## Requirements

- Python 3.11+
- fastmesh >= 0.2
- aiokafka >= 0.11
