Skip to main content

Kafka DI

kafka-di Python confluent-kafka

Kafka DI is a lightweight framework for event-driven Python applications built on Apache Kafka. It provides declarative consumer handlers, dependency injection, and middleware with a small API inspired by FastAPI.

Status: early-stage project. The PyPI package is kafka-di; import it as kafka_di.

Features

  • Declarative Kafka consumer handlers.
  • Dependency injection through Depends.
  • Consumer and producer middleware.
  • Pydantic-friendly typed configuration and serializers.
  • Synchronous and asynchronous message handlers.

Installation

Install with uv:

uv add kafka-di

Or with pip:

pip install kafka-di

Kafka DI requires Python 3.13 or newer.

Quick start

from kafka_di import Kafka
from kafka_di.config import Config

config = Config(
    bootstrap_servers='localhost:9092',
    group_id='orders-service',
)

app = Kafka(configs=config)


@app.consumer.subscribe('orders.created')
def handle_order(message):
    print(f'Received order: {message.value}')


if __name__ == '__main__':
    app.run()

Configuration

Use Config for the common settings, or pass a dictionary with any confluent-kafka settings.

from kafka_di.config import Config, SecurityProtocols

config = Config(
    bootstrap_servers='localhost:9092',
    group_id='orders-service',
)

sasl_config = Config(
    bootstrap_servers='kafka.example.com:9092',
    group_id='orders-service',
    security_protocol=SecurityProtocols.SASL_SSL,
    sasl_username='service-user',
    sasl_password='secret',
    sasl_mechanisms='PLAIN',
)

raw_config = {
    'bootstrap.servers': 'localhost:9092',
    'group.id': 'orders-service',
    'auto.offset.reset': 'earliest',
}

Producing messages

app.producer is created on first use.

app.producer.produce('orders.created', value='{"id": 42}', key='42')
app.producer.flush()

You can also inject the producer into a handler:

from kafka_di import Depends, Producer


@app.consumer.subscribe('orders.received')
def create_invoice(message, producer: Producer = Depends()):
    producer.produce('invoices.requested', value=message.value)

Consuming messages

Register a consumer through app.consumer:

@app.consumer.subscribe('orders.created', 'orders.updated')
def handle_order(message):
    ...

Or register a separate consumer:

from kafka_di import Consumer

inventory_consumer = Consumer()


@inventory_consumer.subscribe('inventory.changed')
def update_inventory(message):
    ...


app.register_consumer(inventory_consumer)

Async handlers are supported:

@app.consumer.subscribe('orders.created')
async def handle_order(message):
    await persist_order(message.value)

Dependency injection

Use Depends to resolve dependencies for a handler:

from kafka_di import Depends


def get_database():
    return DatabaseConnection()


@app.consumer.subscribe('events')
def handle_event(message, database=Depends(get_database)):
    database.save(message.value)

Middleware

Middleware runs before a consumer handler:

from kafka_di.middlewares import ConsumerMiddleware


class LoggingMiddleware(ConsumerMiddleware):
    def handle(self, message):
        print(f'Processing a message from {message.topic}')


app.register_middleware(LoggingMiddleware())

Development

The project uses uv for dependency management.

uv sync --all-groups
make lint
make test

Integration tests start Kafka with Docker and Testcontainers. To run them directly:

make test-integration

For local Kafka and Kafka UI:

docker compose up -d

Kafka is available at localhost:9092 and Kafka UI at http://localhost:8080.

Releasing

Publishing a GitHub Release runs the PyPI workflow. Before the first release, configure a PyPI Trusted Publisher for the GitHub repository and the publish.yml workflow. See CONTRIBUTING.md for the release checklist.

License

Kafka DI is distributed under the MIT License.

Release files for kafka-di 0.1.3

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for kafka-di 0.1.3
File Size Uploaded
kafka_di-0.1.3.tar.gz 42.3 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for kafka-di 0.1.3
File Interpreter ABI Platform
kafka_di-0.1.3-py3-none-any.whl Python 3 none any Details

Total release size: 54.8 kB

Release files / kafka_di-0.1.3.tar.gz

Download URL kafka_di-0.1.3.tar.gz
Size 42.3 kB
Tags Source
SHA-256 checksum
How to use checksums
d14d2ecaa138622c4e3e287829c6b1ff7c337166addfd591932de91a48e38410
BLAKE2b-256 checksum
How to use checksums
d446c043f9d15fe071ee778c9ac9b8ebe4c3ed2d9a555db7db5afde02dcf729d
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 26, 2026.

Transparency log

Release files / kafka_di-0.1.3-py3-none-any.whl

Download URL kafka_di-0.1.3-py3-none-any.whl
Size 12.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
95dd78a72126c3daaf7de77f49e6796af1a7673fe56c0ad0925cc87c90fd3b35
BLAKE2b-256 checksum
How to use checksums
3506f3dff94ec6217347b2a98e4a522a70e134170083e4bed3f648d62c6daecc
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Sep 26, 2026.

Transparency log

Release history Release notifications | RSS feed

This release

0.1.3 This release

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