openframe-adapters-queue-kafka
Apache Kafka queue adapter for the OpenFrame Microservice Development Suite.
Part of the openframe-adapters monorepo.
What it provides
| Symbol | Purpose |
|---|---|
KafkaSettings |
Pydantic-settings subclass — reads all config from env vars |
KafkaProducer[T] |
Generic async message producer — BaseProducer[T] |
KafkaConsumer[T] |
Generic async message consumer — BaseConsumer[T] |
KafkaPlugin |
BasePort (Identity + Lifecycle) — structured lifecycle via PluginRegistry/ApplicationBootstrap |
Installation
# Via meta-package (recommended)
pip install "openframe-adapters[kafka]"
# Or directly
pip install openframe-adapters-queue-kafka
Quick start
from openframe.adapters.queue.kafka import KafkaSettings, KafkaProducer, KafkaConsumer
settings = KafkaSettings(kafka_bootstrap_servers="localhost:9092")
# Produce
producer = KafkaProducer(settings)
await producer.start()
await producer.publish({"event": "item.created", "id": "abc"})
await producer.publish_batch([{"event": "x"}, {"event": "y"}])
await producer.close()
# Consume
consumer = KafkaConsumer(settings)
async def handle(event: dict) -> None:
print(f"Received: {event}")
await consumer.subscribe(handle) # runs until consumer.close() called
Configuration
| Env var | Default | Description |
|---|---|---|
KAFKA_BOOTSTRAP_SERVERS |
required | host:port,host:port |
KAFKA_TOPIC |
"openframe" |
Default topic for producer and consumer |
KAFKA_GROUP_ID |
"openframe-group" |
Consumer group ID |
KAFKA_AUTO_OFFSET_RESET |
"earliest" |
"earliest" or "latest" |
KAFKA_MAX_POLL_RECORDS |
10 |
Max messages per poll |
KAFKA_SESSION_TIMEOUT_MS |
30000 |
Consumer session timeout |
KAFKA_REQUEST_TIMEOUT_MS |
30000 |
Broker request timeout |
KAFKA_SECURITY_PROTOCOL |
"PLAINTEXT" |
"PLAINTEXT", "SSL", "SASL_PLAINTEXT" |
KAFKA_SASL_MECHANISM |
"" |
"PLAIN", "SCRAM-SHA-256", etc. |
KAFKA_SASL_USERNAME |
"" |
SASL username |
KAFKA_SASL_PASSWORD |
"" |
SASL password |
Typed domain objects
from openframe.adapters.queue.kafka import KafkaProducer, KafkaConsumer, KafkaSettings
from dataclasses import dataclass, asdict
@dataclass
class OrderEvent:
order_id: str
event_type: str
class OrderProducer(KafkaProducer[OrderEvent]):
def _serialise(self, message: OrderEvent) -> bytes:
import json
return json.dumps(asdict(message)).encode("utf-8")
class OrderConsumer(KafkaConsumer[OrderEvent]):
def _deserialise(self, raw: bytes) -> OrderEvent:
import json
return OrderEvent(**json.loads(raw.decode("utf-8")))
Plugin lifecycle (optional)
KafkaPlugin is a BasePort — wire it up with ApplicationBootstrap.compose(),
the recommended zero-subclass entry point:
from openframe.core.runtime import ApplicationBootstrap
from openframe.core.ports import Capability
from openframe.adapters.queue.kafka import KafkaPlugin, KafkaSettings
kafka = KafkaPlugin(KafkaSettings())
async with ApplicationBootstrap.compose(kafka) as app:
plugin = app.get(Capability.QUEUE)
producer = plugin.get_producer()
await producer.publish({"event": "item.created"})
consumer = plugin.make_consumer()
await consumer.subscribe(handler)
KafkaPlugin doubles as both the producer and consumer port for the QUEUE
capability (get_producer() / make_consumer()), so a single instance is
usually enough. If a service registers a separate producer-only and
consumer-only KafkaPlugin (e.g. different topics/settings for each), pass
both to compose(): ApplicationBootstrap.compose(producer_plugin, consumer_plugin). Reach for a subclassed ApplicationBootstrap with
configure() only when you need per-port config=/init_timeout= or
conditional registration order, and use app.registry as an escape hatch
for anything neither tier covers.
Consumer acknowledgement semantics
| Outcome | Behaviour |
|---|---|
| Handler returns | ack() called → offset committed → message consumed |
| Handler raises | nack() called → no commit → message redelivered |
consumer.close() |
polling loop exits → consumer stopped cleanly |
License
MIT — © Furious Meteors Engineering
Metadata
Release files for openframe-adapters-queue-kafka 1.4.5
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| openframe_adapters_queue_kafka-1.4.5.tar.gz | 18.3 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| openframe_adapters_queue_kafka-1.4.5-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 32.5 kB
Release files / openframe_adapters_queue_kafka-1.4.5.tar.gz
| Download URL | openframe_adapters_queue_kafka-1.4.5.tar.gz |
|---|---|
| Size | 18.3 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
2650dd8e9d53dba358219f7a6e656c40f1bcdc8583c0309b55e3276354a7e7c8
|
|
BLAKE2b-256 checksum How to use checksums |
08459be37634856c436e8834ab02d03462f818a068d4a93b58211dda3b5dd534
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Release files / openframe_adapters_queue_kafka-1.4.5-py3-none-any.whl
| Download URL | openframe_adapters_queue_kafka-1.4.5-py3-none-any.whl |
|---|---|
| Size | 14.3 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
78c8a584d3a9106b85a4543624eb12b1889af8481962cacf5608c432c6d43bd0
|
|
BLAKE2b-256 checksum How to use checksums |
a6171ef939527f8e77e0570491022fa37f69aff677c983d4f1fbce8e13b1c6bb
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|