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.3.0.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.3.0.tar.gz.

File metadata

File hashes

Hashes for openframe_adapters_queue_kafka-1.3.0.tar.gz
Algorithm Hash digest
SHA256 6cdb81d3813a500581102811f821018dc9e01b8398e6680f692b5ba436089fb2
MD5 e9d79cb30a5e81a329d0e5b62545c523
BLAKE2b-256 6a2d99a4567a960fac5efbe33e7934a4a09782908f33e89a3d16ddb00b50392f

See more details on using hashes here.

File details

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

File metadata

File hashes

Hashes for openframe_adapters_queue_kafka-1.3.0-py3-none-any.whl
Algorithm Hash digest
SHA256 c3f0d0e8ff130207f28c16ef190b8d79deb3242c4b91c8a19fc5644e668aff0a
MD5 72e44ecf2e539b019fe8b212f0f42211
BLAKE2b-256 207bd286bbb13598d7f576a8642678054d01dea0e38eef864dc4b038f40159f3

See more details on using hashes here.

Release history Release notifications | RSS feed

1.4.3

2 files

1.4.2

2 files

1.4.1

2 files

1.4.0

2 files

This release

1.3.0 This release

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