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.4

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.4
File Size Uploaded
openframe_adapters_queue_kafka-1.4.4.tar.gz 17.6 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for openframe-adapters-queue-kafka 1.4.4
File Interpreter ABI Platform
openframe_adapters_queue_kafka-1.4.4-py3-none-any.whl Python 3 none any Details

Total release size: 31.6 kB

Release files / openframe_adapters_queue_kafka-1.4.4.tar.gz

Download URL openframe_adapters_queue_kafka-1.4.4.tar.gz
Size 17.6 kB
Tags Source
SHA-256 checksum
How to use checksums
5e94807295198d457d51ed2a40696b28f8929493ecee107b3dc330da6cb937b7
BLAKE2b-256 checksum
How to use checksums
e0b8ef5504fae73e701268f32cadd8ab61c3791a1b61c4b212b90b3ba3c929bc
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.4-py3-none-any.whl

Download URL openframe_adapters_queue_kafka-1.4.4-py3-none-any.whl
Size 13.9 kB
Tags Python 3
SHA-256 checksum
How to use checksums
ffda461bda6ada72131895536bd6719db289250c921d157d18c765c6c8e9fe24
BLAKE2b-256 checksum
How to use checksums
620feb844e22b57552769bb02f729ee28555606e7f5de6142abdd754a560e599
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

1.4.5

2 release files

This release

1.4.4 This release

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