WKafka v1.0.0 LTS 🚀
Professional, Decorator-based Kafka Wrapper for Python.
WKafka simplifies Apache Kafka integration by providing a high-level, intuitive API focused on developer productivity. It includes built-in support for complex data types like JSON, YAML, Images, Files, and Pydantic models, making it ideal for modern microservices, IoT, and Computer Vision pipelines.
🌟 Features
- Decorator-driven API: Minimalistic and clean message handling.
- Modern Python: Fully typed, PEP 8 compliant, supporting Python 3.9 through 3.14.
- Enterprise Security: Built-in support for SASL (PLAIN, SCRAM) and KRaft mode.
- Multimedia & File Native: Seamlessly send and receive images (OpenCV/NumPy/PIL) and arbitrary files (PDF, ZIP, TXT) via
format="file". - Type-safe Pydantic Validation: Automatic schema validation with
format="pydantic". - Manual Offset Commit: Control At-Least-Once delivery semantics with
auto_commit=Falseandmsg.commit(). - Partition Auto-scaling & Retries: Automatic topic partition scaling with
describe_topicsinspection and exponential backoff retry resilience against transientNodeNotReadyError. - Async/Await Support: Define non-blocking
async defconsumer handlers. - Professional Ops: Structured logging via
loguruand multi-version testing withtox.
📦 Installation
# Via pip
pip install wkafka
# Via poetry
poetry add wkafka
Optional snappy compression:
pip install wkafka[snappy]
🚀 Quick Start
Basic Producer & Consumer
from wkafka import WKafka
# Configures automatically via KAFKA_SERVER or defaults to localhost:9092
kafka = WKafka(bootstrap_servers="localhost:9092")
@kafka.consumer(topic="orders", format="json")
def handle_order(msg):
print(f"New order received: {msg.value}")
# Start consumers in a background thread pool
kafka.run_consumers(block=True)
# Produce with context manager safety
with kafka.producer() as p:
p.send("orders", value={"id": 123, "item": "Coffee"}, format="json")
Manual Offset Commit
@kafka.consumer(topic="transactions", format="json", auto_commit=False)
def handle_tx(msg):
# Process business logic
save_to_db(msg.value)
# Explicitly commit offset only after success
msg.commit()
Retries & DLQ Routing
@kafka.consumer(
topic="unstable_events",
format="json",
max_retries=3,
retry_delay=1.0,
dlq_topic="unstable_events.DLQ"
)
def handle_event(msg):
process_payload(msg.value)
📂 Project Structure
wkafka.core: Orchestration and base logic (WKafka,Message).wkafka.serializers: Extensible serialization system (JSONSerializer,YAMLSerializer,ImageSerializer,PydanticSerializer,FileSerializer).wkafka.controller: Backward compatibility layer for legacy code.examples/: 12 complete, production-ready example modules (01_basicthrough12_pydantic_validation).enviroment/: Production-ready Docker setups (KRaft, SASL).
🛠️ Tecnologías y Librerías Relevantes
- Python (3.9 - 3.14): Lenguaje principal de desarrollo y ejecución.
- kafka-python-ng / kafka-python: Cliente subyacente para comunicación de bajo nivel con Apache Kafka.
- OpenCV (
opencv-python) & Pillow: Procesamiento, renderizado, serialización y deserialización de imágenes. - Pydantic: Validación de esquemas y modelos de datos tipados (
format="pydantic"). - NumPy: Manejo de estructuras de datos matriciales multidimensionales para imágenes.
- PyYAML: Serialización y deserialización nativa de estructuras YAML.
- Loguru: Sistema avanzado de logging estructurado.
🧪 Unit Testing & Coverage
To run the unit test suite, calculate coverage, or execute tests inside a sandboxed Docker container:
Run Locally (pytest)
poetry run pytest tests/
Coverage Report (80% Minimum Guaranteed)
./run_coverage.sh
El script genera un informe de cobertura HTML (htmlcov/index.html) garantizando una cobertura mínima global de 80%.
Sandboxed Testing with Docker
./run_tests_docker.sh
📜 License
MIT License. Created by wisrovi.
Metadata
Release files for wkafka 1.3.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| wkafka-1.3.0.tar.gz | 13.4 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| wkafka-1.3.0-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 26.7 kB
Release files / wkafka-1.3.0.tar.gz
| Download URL | wkafka-1.3.0.tar.gz |
|---|---|
| Size | 13.4 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
ce01e81210d5e13a9dc16b77b26242601a835d6c0e9c2c745074f52a6ec26307
|
|
BLAKE2b-256 checksum How to use checksums |
3cb71178040b971aa9ce1723b3813561e5417ec1c2f45480923ea7f18c4cef19
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.13.5
|
Release files / wkafka-1.3.0-py3-none-any.whl
| Download URL | wkafka-1.3.0-py3-none-any.whl |
|---|---|
| Size | 13.3 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
7c3ffa81dabca61acf3b702b168f5f8040516ac0cbfc032844d94e39595a1daa
|
|
BLAKE2b-256 checksum How to use checksums |
7087af19e390426c4afafde49de089b9b30ef9901d90f005d36c9ae5c3625eb0
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/7.0.0 CPython/3.13.5
|