Skip to main content

DataLoom: Um orquestrador de threads leve e eficiente para dados.

Project description

🧵 DataLoom

Um motor de orquestração multi-thread leve, eficiente e seguro para Python.

CI Python License

DataLoom: Um motor de orquestração multi-thread leve, eficiente e seguro para Python.

O DataLoom é uma biblioteca projetada para processar fluxos de dados utilizando o padrão Produtor-Consumidor com múltiplas threads (Weavers). Ele abstrai a complexidade de filas (Queues), sincronização (Locks) e gerenciamento de ciclo de vida, permitindo que você foque apenas na lógica de transformação dos dados.

Por que usar o DataLoom?

⚡ Performance: Processamento paralelo real para tarefas I/O bound. 🧩 Simplicidade: API intuitiva inspirada na metáfora de tecelagem. 🛡️ Segurança: Thread-safety garantido por design em todo o pipeline. 📦 Leve: Dependências mínimas, pronto para rodar em qualquer ambiente Python.

Por que não usar só ThreadPoolExecutor?

Pergunta justa — a stdlib resolve o paralelismo, mas não o pipeline. O concurrent.futures te dá um pool de threads; tudo ao redor fica por sua conta:

Você precisa de... Com a stdlib Com o DataLoom
Fluxo contínuo produtor-consumidor Queue + loops manuais Loom + Source
Backpressure (produtor mais rápido que consumidores) Queue(maxsize=...) manual Padrão, configurável
Shutdown limpo (drenar fila, fechar recursos) Sentinelas e joins manuais stop() / with Loom(...)
Worker que sobrevive a erros e os reporta try/except em cada worker hooks.on_error centralizado
Métricas por item processado Instrumentação manual hooks.on_batch_processed
Escrita concorrente segura em arquivo Lock manual Sinks prontos (JSON, CSV, callback)

Se o seu caso é "aplicar uma função a uma lista e coletar os resultados", use ThreadPoolExecutor.map — é a ferramenta certa. O DataLoom é para fluxos contínuos ou longos onde ciclo de vida, resiliência e observabilidade importam. E, como toda solução baseada em threads no CPython, o ganho de paralelismo vale para cargas I/O bound; para CPU bound, prefira multiprocessing.

✨ Características

  • API Coesa: Conceitos alinhados (Loom, Weaver, Sink).
  • Concorrente: Gerencia automaticamente múltiplos Weavers (threads) em paralelo.
  • Seguro: Sinks padrão são thread-safe e exceções são tipadas (LoomError).
  • Resiliente: Erros de processamento não derrubam os Weavers e são reportados via hooks.
  • Gerenciável: Loom é um context manager — with Loom(...) as loom: garante stop() automático.
  • Backpressure: Fila de tarefas limitada por padrão (queue_maxsize), evitando crescimento de memória sem controle.
  • Observável: Hooks de ciclo de vida e métricas por lote (on_batch_processed com resultado e duração), além de Logs integrados.

📦 Instalação

pip install dataloom-engine

⚠️ Atenção ao nome: o pacote instala o módulo dataloom_engine (import dataloom_engine). Não confunda com o pacote dataloom do PyPI, que é um ORM de outro autor e não tem relação com este projeto.

Para desenvolver ou usar a versão mais recente do repositório:

git clone https://github.com/dionipadilha/dataloom.git
cd dataloom
pip install -e .

🚀 Uso Rápido

Aqui está um exemplo completo de como tecer um pipeline de dados:

from pathlib import Path
import numpy as np
from dataloom_engine import (
    Loom,
    LoomConfig,
    LoomLogs,
    Processor,
    JsonFileSink,
    ThreadedBufferedSink
)
from dataloom_engine.sources import RandomNumPySource

# 1. Defina sua lógica de processamento (Stateless)
class MyFilterProcessor(Processor):
    def process(self, batch: np.ndarray) -> dict:
        # 'batch' é um array numpy com o tamanho definido na config
        avg = float(batch.mean())
        return {
            "processed_items": len(batch),
            "average_value": avg,
            "status": "high" if avg > 0.5 else "low"
        }

# 2. Configuração Inicial
if __name__ == "__main__":
    # Configura logs no console
    LoomLogs.setup()

    # Define parâmetros do motor
    config = LoomConfig(
        output_dir=Path("./data_out"),
        batch_size=100,      # Processa 100 itens por vez
        interval_seconds=1   # Gera um novo lote a cada 1 segundo
    )

    # 3. Fonte de Dados
    # Agora desvinculada do motor, permitindo fontes customizadas (DB, CSV, S3)
    source = RandomNumPySource(config)

    # 4. Sink Arquivado e Otimizado
    # ThreadedBufferedSink evita que o I/O bloqueie o processamento
    file_sink = JsonFileSink(config.output_dir)
    sink = ThreadedBufferedSink(file_sink)

    # 5. Inicializa o Loom (O Orquestrador)
    loom = Loom(
        config=config,
        processor=MyFilterProcessor(),
        sink=sink,
        source=source,
        num_weavers=4  # 4 Threads trabalhando em paralelo
    )

    print("🧵 DataLoom iniciado! Pressione Ctrl+C para parar.")
    try:
        # O context manager garante stop() e limpeza dos recursos,
        # mesmo em caso de exceção ou Ctrl+C
        with loom:
            loom.start()
    except KeyboardInterrupt:
        print("\n🛑 Tear parado.")

🏗️ Arquitetura

O DataLoom utiliza uma metáfora de tecelagem:

  • Loom (Tear): A máquina principal. Gerencia a fila de tarefas e o ciclo de vida.
  • Weaver (Tecelão): As threads trabalhadoras. Elas pegam a matéria-prima (batch), processam e entregam.
  • Processor: A lógica de negócio. Transforma dados brutos em informação.
  • Sink: O destino final. Onde o produto acabado é depositado (ex: JsonFileSink, CsvFileSink, ou qualquer destino via CallbackSink).

🛠️ Desenvolvimento e Testes

Para contribuir com o projeto ou rodar os testes unitários:

  1. Instale as dependências de desenvolvimento:

    pip install -e ".[dev]"
    
  2. Rode a suíte de testes (via pytest):

    pytest
    

📄 Licença

Este projeto está licenciado sob a licença MIT - veja o arquivo LICENSE para detalhes.

🗒️ Histórico de versões

As mudanças de cada versão estão documentadas no CHANGELOG.

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

dataloom_engine-0.3.0.tar.gz (22.2 kB view details)

Uploaded Source

Built Distribution

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

dataloom_engine-0.3.0-py3-none-any.whl (16.8 kB view details)

Uploaded Python 3

File details

Details for the file dataloom_engine-0.3.0.tar.gz.

File metadata

  • Download URL: dataloom_engine-0.3.0.tar.gz
  • Upload date:
  • Size: 22.2 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/6.1.0 CPython/3.13.12

File hashes

Hashes for dataloom_engine-0.3.0.tar.gz
Algorithm Hash digest
SHA256 67cf163ee6abc2378ca833d755858f56eb2462cb65f2be1fa19838498d67d916
MD5 c7d5ff7aa05dd9aa17af3da9d3d2bc56
BLAKE2b-256 6068ec47717df858095514748262f9ad836ee3b8bcc3c3753477e0a461a38aba

See more details on using hashes here.

Provenance

The following attestation bundles were made for dataloom_engine-0.3.0.tar.gz:

Publisher: publish.yml on dionipadilha/dataloom

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file dataloom_engine-0.3.0-py3-none-any.whl.

File metadata

  • Download URL: dataloom_engine-0.3.0-py3-none-any.whl
  • Upload date:
  • Size: 16.8 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/6.1.0 CPython/3.13.12

File hashes

Hashes for dataloom_engine-0.3.0-py3-none-any.whl
Algorithm Hash digest
SHA256 33f6d6e2db0661d90e17c1645269da7bc3e1a52f35a7d8d22b60614d469ae8a1
MD5 8a63ded4aaa282786f5841f21b6a5b54
BLAKE2b-256 265e9a565f064a8b5f24128fb360bc54d359ba514f47d54893d0d30767069da5

See more details on using hashes here.

Provenance

The following attestation bundles were made for dataloom_engine-0.3.0-py3-none-any.whl:

Publisher: publish.yml on dionipadilha/dataloom

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

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