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.
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:garantestop()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_processedcom 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 pacotedataloomdo 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 viaCallbackSink).
🛠️ Desenvolvimento e Testes
Para contribuir com o projeto ou rodar os testes unitários:
-
Instale as dependências de desenvolvimento:
pip install -e ".[dev]"
-
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
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
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
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
67cf163ee6abc2378ca833d755858f56eb2462cb65f2be1fa19838498d67d916
|
|
| MD5 |
c7d5ff7aa05dd9aa17af3da9d3d2bc56
|
|
| BLAKE2b-256 |
6068ec47717df858095514748262f9ad836ee3b8bcc3c3753477e0a461a38aba
|
Provenance
The following attestation bundles were made for dataloom_engine-0.3.0.tar.gz:
Publisher:
publish.yml on dionipadilha/dataloom
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
dataloom_engine-0.3.0.tar.gz -
Subject digest:
67cf163ee6abc2378ca833d755858f56eb2462cb65f2be1fa19838498d67d916 - Sigstore transparency entry: 2154260466
- Sigstore integration time:
-
Permalink:
dionipadilha/dataloom@c21e3bc70bbcc47c19454c37b43aadcedfc97d2b -
Branch / Tag:
refs/tags/v0.3.0 - Owner: https://github.com/dionipadilha
-
Access:
private
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
publish.yml@c21e3bc70bbcc47c19454c37b43aadcedfc97d2b -
Trigger Event:
release
-
Statement type:
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
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
33f6d6e2db0661d90e17c1645269da7bc3e1a52f35a7d8d22b60614d469ae8a1
|
|
| MD5 |
8a63ded4aaa282786f5841f21b6a5b54
|
|
| BLAKE2b-256 |
265e9a565f064a8b5f24128fb360bc54d359ba514f47d54893d0d30767069da5
|
Provenance
The following attestation bundles were made for dataloom_engine-0.3.0-py3-none-any.whl:
Publisher:
publish.yml on dionipadilha/dataloom
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
dataloom_engine-0.3.0-py3-none-any.whl -
Subject digest:
33f6d6e2db0661d90e17c1645269da7bc3e1a52f35a7d8d22b60614d469ae8a1 - Sigstore transparency entry: 2154260538
- Sigstore integration time:
-
Permalink:
dionipadilha/dataloom@c21e3bc70bbcc47c19454c37b43aadcedfc97d2b -
Branch / Tag:
refs/tags/v0.3.0 - Owner: https://github.com/dionipadilha
-
Access:
private
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
publish.yml@c21e3bc70bbcc47c19454c37b43aadcedfc97d2b -
Trigger Event:
release
-
Statement type: