Skip to main content

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)

Source distribution for openframe-adapters-queue-kafka 1.4.5
File Size Uploaded
openframe_adapters_queue_kafka-1.4.5.tar.gz 18.3 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for openframe-adapters-queue-kafka 1.4.5
File Interpreter ABI Platform
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

Release history Release notifications | RSS feed

This release

1.4.5 This release

2 release files

1.4.4

2 release files

1.4.3

2 release files

1.4.2

2 release files

1.4.1

2 release files

1.4.0

2 release files

1.3.0

2 release files

1.2.0

2 release files

1.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page