Skip to main content
image

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=False and msg.commit().
  • Partition Auto-scaling & Retries: Automatic topic partition scaling with describe_topics inspection and exponential backoff retry resilience against transient NodeNotReadyError.
  • Async/Await Support: Define non-blocking async def consumer handlers.
  • Professional Ops: Structured logging via loguru and multi-version testing with tox.

📦 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_basic through 12_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)

Source distribution for wkafka 1.3.0
File Size Uploaded
wkafka-1.3.0.tar.gz 13.4 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for wkafka 1.3.0
File Interpreter ABI Platform
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

Release history Release notifications | RSS feed

This release

1.3.0 This release

2 release files

1.2.2

2 release files

1.2.1

2 release files

1.2.0

2 release files

1.1.9

2 release files

1.0.9

2 release files

1.0.7

2 release files

1.0.6

2 release files

1.0.5

2 release files

1.0.4

2 release files

1.0.3

2 release files

1.0.2

2 release files

1.0.1

2 release files

1.0.0

2 release files

0.2.1

2 release files

0.2.0

2 release files

0.1.9

2 release files

0.1.8

2 release files

0.1.7

2 release files

0.1.6

2 release files

0.1.5

2 release files

0.1.4

2 release files

0.1.3

2 release files

0.1.2

2 release files

0.1.1

2 release files

0.0.1

1 release file

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