Kafka DI
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 askafka_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)
| File | Size | Uploaded | |
|---|---|---|---|
| kafka_di-0.1.3.tar.gz | 42.3 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| 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 logRelease 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