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 OpenFramePlugin — structured lifecycle via PluginRegistry

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)

from openframe.core.plugins import PluginRegistry
from openframe.adapters.queue.kafka import KafkaPlugin, KafkaSettings

registry = PluginRegistry()
registry.register(KafkaPlugin(KafkaSettings()))
await registry.initialize_all()

plugin = registry.get("queue")
producer = plugin.get_producer()
await producer.publish({"event": "item.created"})

consumer = plugin.make_consumer()
await consumer.subscribe(handler)

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

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

openframe_adapters_queue_kafka-1.4.2.tar.gz (16.6 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

File details

Details for the file openframe_adapters_queue_kafka-1.4.2.tar.gz.

File metadata

File hashes

Hashes for openframe_adapters_queue_kafka-1.4.2.tar.gz
Algorithm Hash digest
SHA256 8d6273b0b2c8f0e6cf2551cbc8d239eb1ca74070fceda3618597ee97eca2998b
MD5 55aa0ecf59bff3e9cfca1525c90fe00d
BLAKE2b-256 fb6882db0c45a5c3c0b9c9743015005caa4923aef3c50aa29757a30f2a819dea

See more details on using hashes here.

File details

Details for the file openframe_adapters_queue_kafka-1.4.2-py3-none-any.whl.

File metadata

File hashes

Hashes for openframe_adapters_queue_kafka-1.4.2-py3-none-any.whl
Algorithm Hash digest
SHA256 907010d7cdcee80ef289cbb4f28af60b2eb56bb9e5aa3553d5396cebb2dcf4ba
MD5 7c719f5c4f4f2b839c8539a76c5ef09b
BLAKE2b-256 b40ff2111cda9f42ea480926f1d9ec6f4badf62d196ff39f609a3cf5f5db393e

See more details on using hashes here.

Release history Release notifications | RSS feed

1.4.3

2 files

This release

1.4.2 This release

2 files

1.4.1

2 files

1.4.0

2 files

1.3.0

2 files

1.2.0

2 files

1.1.0

2 files

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page