Skip to main content

Kafka Integration for Muffin framework

Project description

Muffin-Kafka

Muffin-Kafka is an Apache Kafka integration plugin for the Muffin web framework, built on top of aiokafka.

Tests Status PYPI Version Python Versions


🚀 Features

  • Async Kafka integration using aiokafka
  • Per-topic task model — each topic is consumed in an isolated asyncio task
  • Simple handler registration using @plugin.handle_topics(...)
  • Manual or auto-commit support, custom group IDs
  • Producer support (send / send_and_wait)
  • Built-in monitoring with offsets, lag, and poll delay
  • Healthcheck support for liveness probes and observability
  • Optional error handler via @plugin.handle_error(...)

✅ Requirements

  • Python ≥ 3.10
  • Muffin ≥ 0.71
  • Kafka cluster or broker (local or cloud)

📦 Installation

pip install muffin-kafka

⚙️ Usage

    from muffin import Application
    from muffin_kafka import Kafka

    app = Application("example")

    # Initialize the plugin with config options
    kafka = Kafka(app, bootstrap_servers="localhost:9092", produce=True, listen=True)

🧩 Registering Handlers

Use @kafka.handle_topics(...) to register a handler for specific Kafka topics:

    @kafka.handle_topics("events.user", "events.auth")
    async def handle_event(message):
        data = message.value.decode()
        print("Received:", data)

You can also register a global error handler:

    @kafka.handle_error
    async def on_error(exc):
        print("Kafka error:", exc)

📤 Sending Messages

You can send messages to Kafka topics using the send or send_and_wait methods:

    # Send a message without waiting for acknowledgment
    await kafka.send("events.user", {"action": "signup"}, key="user123")

    # Or wait for broker acknowledgment
    result = await kafka.send_and_wait("events.user", {"action": "login"})

🔄 Healthcheck

You can monitor consumer health by checking lag across partitions:

    # Check if Kafka lag is within acceptable limits
    ok = await kafka.healthcheck(max_lag=1000)
    if not ok:
        raise RuntimeError("Kafka lag too high")

📊 Monitoring

If monitor=True is passed, the plugin will log:

  • Committed offsets
  • Latest offsets
  • Poll timestamps
  • Per-partition lag and delay

This data can be extended for Prometheus/Grafana metrics or alerting.

⚙️ Configuration Options

You can pass configuration options either as keyword arguments to the plugin:

kafka = Kafka(app, bootstrap_servers="localhost:9092", produce=True)

Or set them via Muffin's config system (e.g. .env, YAML):

"KAFKA_BOOTSTRAP_SERVERS": "localhost:9092",
"KAFKA_PRODUCE": True,

Supported Options

Option Type Default Description
bootstrap_servers str "localhost:9092" Kafka broker connection string
group_id str None Kafka consumer group ID
client_id str "muffin" Kafka client ID
produce bool False Enable Kafka producer
listen bool True Enable consumers (message listening)
monitor bool False Enable internal consumer monitor
monitor_interval int 60 Monitor frequency in seconds
auto_offset_reset str "earliest" Where to start if no committed offset
enable_auto_commit bool False Automatically commit offsets
max_poll_records int None Max records to poll in one batch
request_timeout_ms int 30000 Request timeout
retry_backoff_ms int 1000 Retry interval on failure
security_protocol str "PLAINTEXT" Kafka protocol (SSL, SASL_PLAINTEXT, …)
sasl_mechanism str "PLAIN" SASL auth mechanism
sasl_plain_username str None SASL auth user
sasl_plain_password str None SASL auth password
ssl_cafile str None Path to trusted CA certs

🐞 Bug Tracker

Found a bug or have a feature request? Please open an issue at: https://github.com/klen/muffin-kafka/issues


🤝 Contributing

Pull requests are welcome! Development happens here: https://github.com/klen/muffin-kafka


🪪 License

MIT – See LICENSE for full details.

Project details


Download files

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

Source Distribution

muffin_kafka-1.0.3.tar.gz (6.8 kB view details)

Uploaded Source

Built Distribution

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

muffin_kafka-1.0.3-py3-none-any.whl (7.7 kB view details)

Uploaded Python 3

File details

Details for the file muffin_kafka-1.0.3.tar.gz.

File metadata

  • Download URL: muffin_kafka-1.0.3.tar.gz
  • Upload date:
  • Size: 6.8 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: uv/0.11.1 {"installer":{"name":"uv","version":"0.11.1","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}

File hashes

Hashes for muffin_kafka-1.0.3.tar.gz
Algorithm Hash digest
SHA256 65483ae0294ce546fba60d79c7981f045fca58ca217816a17624b89ba9c5c045
MD5 9a9ac9a3df358b2b20b080f6858073d1
BLAKE2b-256 db87fc25e10597451ee0c306516a57853f036216762ccfec42d0bb06cf007dab

See more details on using hashes here.

File details

Details for the file muffin_kafka-1.0.3-py3-none-any.whl.

File metadata

  • Download URL: muffin_kafka-1.0.3-py3-none-any.whl
  • Upload date:
  • Size: 7.7 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: uv/0.11.1 {"installer":{"name":"uv","version":"0.11.1","subcommand":["publish"]},"python":null,"implementation":{"name":null,"version":null},"distro":{"name":"Ubuntu","version":"24.04","id":"noble","libc":null},"system":{"name":null,"release":null},"cpu":null,"openssl_version":null,"setuptools_version":null,"rustc_version":null,"ci":true}

File hashes

Hashes for muffin_kafka-1.0.3-py3-none-any.whl
Algorithm Hash digest
SHA256 23e415f96f749d0acf0a91dab5c45ca6a8b5678560f4247f551c9764717cb638
MD5 d2c5c74c62fc5e8be0bcf903c02930b7
BLAKE2b-256 f053c4ce6bb6e1b0eba1d5e928cc98d164a071239f0b70576ef642be1c92f462

See more details on using hashes here.

Supported by

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