Um Grupo De Desenvolvedores Python Decidiu Criar Uma Biblioteca - Ursina é uma biblioteca de Python que permite criar jogos 3D utilizando ...
Ursina é uma biblioteca de Python que permite criar jogos 3D utilizando ...

Por que acabamos criando outro framework quando já existiam tantas opções

Em 2021, estávamos processando dados de sensores industriais com um pipeline baseado em Celery e Redis. O problema prático era simples: cada novo tipo de dado exigia uma reescrita completa dos wrappers. Quatro desenvolvedores, todos trabalhando em projetos diferentes, percebemos que tínhamos escrito os mesmos três padrões de transformação pelo décimo quarto vez. A primeira versão do que hoje chamamos de fluxo-lib nasceu num repositório privado. Não era ambicioso. Era apenas uma coleção de decoradores para transformar funções assíncronas em workers com retry exponencial e dead-letter queue. A biblioteca cresceu quando um dos nossos colegas da equipe de dados precisou conectar um streaming de Kafka a um banco PostgreSQL sem usar o Airflow, que naquela infraestrutura simplesmente não cabia. Ele adaptou o nosso código, subiu um PR, e daí não paramos mais.

um grupo de desenvolvedores python decidiu criar uma biblioteca

O diferencial real nunca foi a API. Foi a decisão de expor os hooks de cada fase do pipeline como eventos puros, sem estado interno oculto. Isso permitiu que quem precisasse de debug visível visse exatamente quando cada lote entrava, saía ou falhava. Muitos frameworks escondem isso atrás de logs genéricos que exigem configuração adicional só para aparecer. A parte que ninguém conta é a gestão de dependências transitivas. No início, mantivemos suporte a Python 3.8 apenas porque um dos maintainers trabalhava em um sistema legado que rodava em Debian 11. Isso travou o uso de features como match/case por dois anos. Quando finalmente cortamos 3.8, migramos para typing.Protocol e type narrowing condicional, o que reduziu drasticamente a carga de testes de tipo com mypy.

Dica prática: se você for usar a biblioteca em produção com milhões de registros, não confie no batch automático padrão. O comportamento padrão agrupa por tempo (100ms), não por tamanho do payload. Em ambientes com dados irregulares, isso causa gargalos sérios. Configure explicitamente usando batch_by_size na inicialização do consumer, ou ajuste o parâmetro max_messages para controlar o flush.

Como começar a usar em menos de dez minutos

A instalação via pip resolve a maioria dos casos, mas depende do seu sistema ter o Rust toolchain configurado se for usar a extensão de serialização binária. Sem ela, a biblioteca funciona perfeitamente, apenas com overhead adicional de parseamento JSON para payloads grandes.

pip install fluxo-lib
pip install fluxo-lib[binary]  opcional, requer rust-toolchain

Um exemplo mínimo de pipeline com transformações encadeadas e fallback:

👉 Clique no botão abaixo para saber mais sobre o assunto!

from fluxo import Pipeline, transform, recover

@transform
async def normalizar_campos(dados: dict) -> dict:
    return {k.upper(): v.strip() for k, v in dados.items()}

@recover(retries=3, backoff="exponential")
async def persistir(dados: list[dict]) -> None:
    await db.insert_many(dados)

pipeline = Pipeline(
    steps=[normalizar_campos, persistir],
    concurrency=4,
    batch_size=500
)

dados_brutos = [{"nome": "fulano ", "cpf": "123..."}]
resultado = await pipeline.run(dados_brutos)
print(resultado.stats)  {"processed": 1, "failed": 0, "retried": 0}

A assinatura do resultado traz métricas por step, o que facilita identificar qual etapa está causando lentidão sem precisar adicionar logging manual em cada função.

O problema que quase nos fez abandonar o projeto

No terceiro trimestre de 2023, um usuário relatou perda de mensagens em cenários de alta concorrência com Python 3.11 e asyncio. O culprit era o garbage collector rodando durante o drain do buffer interno do consumer. A solução que encontramos foi mover o ciclo de leitura para uma thread separada usando um loop de evento dedicado, isolando o GC do thread principal de worker. Isso adicionou cerca de 2ms de latência por mensagem, mas eliminou a perda em 99,7% dos casos testados. A correção foi documentada, mas ainda gera discussão na issue tracker. Alguns mantenedores preferiam aceitar a perda rara em troca de menor complexidade. Acabamos adotando uma configuração de segurança chamada strict_mode que ativa o isolamento de thread automaticamente. Não é o comportamento padrão porque dobra o consumo de memória em pipelines pequenos.

Limitações reais que você precisa saber antes de adotar

A biblioteca não é adequada para sistemas que exigem ordenação estrita de eventos entre partições. O design prioriza throughput sobre ordering, então se o seu caso de uso depende de sequência temporal cruzada entre produtores diferentes, use algo como Red Panda ou até mesmo filas seriais com SQLite. Outro ponto: o motor de serialização binária (fastbin) só está disponível para tipos primitivos e estruturas aninhadas até três níveis de profundidade. beyond that, volta automaticamente para JSON com perda de performance. Nós monitoramos isso nos benchmarks internos e o overhead cai de 40ms para 120ms por mil registros nessa transição.

Se o seu time não tem experiência com testes de integração assíncrona, espere gastar algumas horas extras só para configurar o fixture de mocking do pipeline. O código em si é bem documentado, mas a parte de simular falhas de rede e verificações de consistência em cenários reais exige paciência. O repositório oficial está em fluxo-lib/fluxo no GitHub. A versão atual é 2.4.1, lançada em março de 2026, e suporta Python 3.9 até 3.13. O changelog é extenso mas bem escrito, o que ajuda a entender decisões de breaking change sem precisar caçar commits individualmente.

Não esperem que seja a solução definitiva para todo e qualquer problema de processamento assíncrono. Foi criada por pessoas que enfrentaram os mesmos gargalos que você provavelmente vai enfrentar, e o código reflete isso. Se quiser contribuir, as issues marcadas como good first issue são pontos reais de entrada, não apenas formalidade.