Pular para o conteúdo

    Capítulo 65, Projetos

    Projeto Avançado: processador concorrente de dados

    Um pipeline que lê muitos arquivos com concorrência limitada, calcula em vários processos, tolera falhas e prova isso com testes. Faça depois do capítulo 51.

    Os arquivos deste capítulo estão em projetos/processador/.

    O problema

    Uma pasta com muitos arquivos CSV de vendas (regiao,valor_centavos). Você precisa somar tudo por região. Cada arquivo pode ser grande, alguns podem estar corrompidos e o disco pode falhar de vez em quando. O programa deve ser rápido, limitar quantos arquivos lê ao mesmo tempo e nunca deixar um arquivo ruim derrubar os outros.

    RequisitoDecisãoCapítulo
    Ler muitos arquivos sem bloquearasyncio com to_thread47
    Limitar leituras simultâneasSemaphore47
    Calcular em paralelo de verdadeProcessPoolExecutor46
    Tolerar falha de leituraRetentativa com espera crescente51
    Não esperar para sempreasyncio.timeout47
    Uma falha não derruba as outrasResultado ou Falha, nunca exceção solta23
    Rastrear cada arquivo nos logsContextVar e filtro de log51
    Código testávelInjeção da espera e da leitura51

    O desenho

    O gargalo tem dois tipos, e cada um recebe a ferramenta certa: ler é espera (disco), então uma thread por arquivo, orquestradas por asyncio. Calcular é CPU, então processos. O asyncio coordena as duas etapas em uma única thread, e não há trava nem estado compartilhado.

    O caminho de cada arquivo
    CSV  ->  [semáforo]  ->  ler (thread, com retentativas e timeout)
                                  |
                                  v
                        agregar (outro processo)  ->  Resultado
                                  |
                             erro em qualquer etapa  ->  Falha
    

    Uma regra de desenho vale mais do que o código: cada arquivo vira um Resultado ou uma Falha, e a função que o processa captura os erros esperados. Assim o asyncio.gather nunca vê uma exceção, e um arquivo corrompido vira uma linha no relatório, e não um programa interrompido.

    Os modelos

    Os dados que cruzam as etapas são dataclasses imutáveis. O Relatorio soma por região:

    projetos/processador/src/processador/modelos.py
    from dataclasses import dataclass, field
    
    
    @dataclass(frozen=True)
    class Resultado:
        arquivo: str
        linhas: int
        total_centavos: int
        por_regiao: dict[str, int]
    
    
    @dataclass(frozen=True)
    class Falha:
        arquivo: str
        motivo: str
    
    
    @dataclass
    class Relatorio:
        resultados: list[Resultado] = field(default_factory=list)
        falhas: list[Falha] = field(default_factory=list)
    
        @property
        def total_centavos(self) -> int:
            return sum(r.total_centavos for r in self.resultados)
    
        @property
        def por_regiao(self) -> dict[str, int]:
            totais: dict[str, int] = {}
            for resultado in self.resultados:
                for regiao, valor in resultado.por_regiao.items():
                    totais[regiao] = totais.get(regiao, 0) + valor
            return totais
    

    O cálculo (a etapa de CPU)

    agregar é uma função pura: recebe texto, devolve números, e levanta ValueError com o número da linha quando algo está errado. Ela precisa estar no nível do módulo e usar argumentos e retorno serializáveis, porque é enviada a outro processo:

    projetos/processador/src/processador/calculo.py
    import csv
    import io
    
    
    def agregar(texto: str) -> tuple[int, int, dict[str, int]]:
        """Etapa de CPU: roda em outro processo. Devolve (linhas, total, total por região).
    
        Precisa ser uma função do nível do módulo, com argumentos e retorno serializáveis,
        para poder ser enviada a um processo filho.
        """
        leitor = csv.reader(io.StringIO(texto))
        cabecalho = next(leitor, None)
        if cabecalho != ["regiao", "valor_centavos"]:
            raise ValueError(f"cabeçalho inesperado: {cabecalho}")
        linhas = total = 0
        por_regiao: dict[str, int] = {}
        for numero, linha in enumerate(leitor, start=2):
            try:
                regiao, valor_texto = linha
                valor = int(valor_texto)
            except ValueError as erro:
                raise ValueError(f"linha {numero} inválida: {linha}") from erro
            por_regiao[regiao] = por_regiao.get(regiao, 0) + valor
            total += valor
            linhas += 1
        return linhas, total, por_regiao
    

    A leitura (a etapa de espera)

    A leitura bloqueante vai para uma thread com asyncio.to_thread, para não parar o laço de eventos. As retentativas aumentam a espera a cada falha, e a função de espera e a de leitura são parâmetros, o que permite aos testes simularem falhas sem esperar de verdade:

    projetos/processador/src/processador/leitura.py
    import asyncio
    import logging
    from collections.abc import Awaitable, Callable
    from pathlib import Path
    
    log = logging.getLogger(__name__)
    
    
    def ler_texto(caminho: Path) -> str:
        return caminho.read_text(encoding="utf-8")
    
    
    async def ler_com_retentativas(
        caminho: Path,
        *,
        tentativas: int = 3,
        base: float = 0.1,
        dormir: Callable[[float], Awaitable[None]] = asyncio.sleep,
        ler: Callable[[Path], str] = ler_texto,
    ) -> str:
        """Lê o arquivo em uma thread, para não bloquear o laço de eventos, com espera crescente."""
        for numero in range(1, tentativas + 1):
            try:
                return await asyncio.to_thread(ler, caminho)
            except OSError as erro:
                if numero == tentativas:
                    raise
                espera = base * 2 ** (numero - 1)
                log.warning(
                    "leitura falhou (%s), tentativa %d, nova tentativa em %.2fs", erro, numero, espera
                )
                await dormir(espera)
        raise AssertionError("inalcançável")
    

    Contexto de log

    Cada arquivo é processado em uma tarefa própria. Uma ContextVar guarda o nome do arquivo em andamento, e um filtro o acrescenta a toda linha de log daquela tarefa, sem passar o nome por parâmetro em cada chamada:

    projetos/processador/src/processador/contexto.py
    import logging
    from contextvars import ContextVar
    
    id_arquivo: ContextVar[str] = ContextVar("id_arquivo", default="-")
    
    
    class FiltroArquivo(logging.Filter):
        """Acrescenta o nome do arquivo em processamento a cada linha de log."""
    
        def filter(self, record: logging.LogRecord) -> bool:
            record.arquivo = id_arquivo.get()
            return True
    

    O pipeline

    Aqui está o coração. Repare em três coisas. A função criar_pool escolhe o método de início spawn de forma explícita: com fork, o processo filho herda as threads do asyncio.to_thread em estado inconsistente (o Python 3.12 já emite um aviso sobre isso, e o padrão muda conforme o sistema e a versão), e o spawn se comporta igual em todos. O Semaphore limita as leituras e o timeout limita o tempo por arquivo. E o except converte cada falha esperada em uma Falha:

    projetos/processador/src/processador/pipeline.py
    import asyncio
    import logging
    import multiprocessing
    from collections.abc import Awaitable, Callable
    from concurrent.futures import Executor, ProcessPoolExecutor
    from pathlib import Path
    
    from processador.calculo import agregar
    from processador.contexto import id_arquivo
    from processador.leitura import ler_com_retentativas, ler_texto
    from processador.modelos import Falha, Relatorio, Resultado
    
    log = logging.getLogger(__name__)
    
    
    def criar_pool(processos: int) -> ProcessPoolExecutor:
        """Pool de processos com o método de início 'spawn' escolhido de forma explícita.
    
        O fork copiaria para o filho threads em estado inconsistente (como as do asyncio.to_thread),
        e o método padrão muda conforme o sistema e a versão do Python. O spawn se comporta igual
        em todos, ao custo de iniciar cada processo importando o código de novo.
        """
        return ProcessPoolExecutor(
            max_workers=processos, mp_context=multiprocessing.get_context("spawn")
        )
    
    
    async def processar_arquivo(
        caminho: Path,
        *,
        semaforo: asyncio.Semaphore,
        pool: Executor,
        timeout: float,
        dormir: Callable[[float], Awaitable[None]],
        ler: Callable[[Path], str],
    ) -> Resultado | Falha:
        token = id_arquivo.set(caminho.name)
        try:
            async with semaforo:
                async with asyncio.timeout(timeout):
                    texto = await ler_com_retentativas(caminho, dormir=dormir, ler=ler)
            laco = asyncio.get_running_loop()
            linhas, total, por_regiao = await laco.run_in_executor(pool, agregar, texto)
            log.info("processado: %d linhas", linhas)
            return Resultado(caminho.name, linhas, total, por_regiao)
        except (OSError, TimeoutError, ValueError) as erro:
            motivo = str(erro) or type(erro).__name__
            log.error("falhou: %s", motivo)
            return Falha(caminho.name, motivo)
        finally:
            id_arquivo.reset(token)
    
    
    async def processar(
        pasta: Path,
        *,
        pool: Executor,
        leituras: int = 4,
        timeout: float = 5.0,
        dormir: Callable[[float], Awaitable[None]] = asyncio.sleep,
        ler: Callable[[Path], str] = ler_texto,
    ) -> Relatorio:
        """Lê vários arquivos com concorrência limitada e calcula em processos separados.
    
        Uma falha em um arquivo vira um registro de Falha e não derruba os demais.
        """
        semaforo = asyncio.Semaphore(leituras)
        saidas = await asyncio.gather(
            *(
                processar_arquivo(
                    caminho, semaforo=semaforo, pool=pool, timeout=timeout, dormir=dormir, ler=ler
                )
                for caminho in sorted(pasta.glob("*.csv"))
            )
        )
        relatorio = Relatorio()
        for saida in saidas:
            if isinstance(saida, Resultado):
                relatorio.resultados.append(saida)
            else:
                relatorio.falhas.append(saida)
        return relatorio
    

    A linha de comando

    Os dados de exemplo são determinísticos (sem aleatoriedade), então o resultado é o mesmo em qualquer máquina e dá para conferir. O código de saída é 1 se algum arquivo falhou, o que permite usar o programa em scripts e no CI:

    projetos/processador/src/processador/cli.py
    import argparse
    import asyncio
    import csv
    import logging
    import sys
    from pathlib import Path
    
    from processador.contexto import FiltroArquivo
    from processador.modelos import Relatorio
    from processador.pipeline import criar_pool, processar
    
    REGIOES = ["norte", "nordeste", "centro-oeste", "sudeste", "sul"]
    
    
    def gerar_dados(pasta: Path, arquivos: int, linhas: int) -> None:
        """Gera arquivos CSV determinísticos (sem aleatoriedade, para o resultado ser repetível)."""
        pasta.mkdir(parents=True, exist_ok=True)
        for i in range(arquivos):
            with (pasta / f"vendas_{i:02d}.csv").open("w", encoding="utf-8", newline="") as arquivo:
                escritor = csv.writer(arquivo)
                escritor.writerow(["regiao", "valor_centavos"])
                for j in range(linhas):
                    escritor.writerow([REGIOES[(i + j) % 5], (i * 7919 + j * 104729) % 50_000 + 100])
    
    
    def formatar_reais(centavos: int) -> str:
        reais, resto = divmod(centavos, 100)
        return f"R$ {reais:,}".replace(",", ".") + f",{resto:02d}"
    
    
    def imprimir(relatorio: Relatorio) -> None:
        print(f"arquivos processados: {len(relatorio.resultados)}  falhas: {len(relatorio.falhas)}")
        for regiao, total in sorted(relatorio.por_regiao.items()):
            print(f"  {regiao:<13}{formatar_reais(total):>16}")
        print(f"  {'total':<13}{formatar_reais(relatorio.total_centavos):>16}")
        for falha in relatorio.falhas:
            print(f"  FALHA {falha.arquivo}: {falha.motivo}")
    
    
    def configurar_logs(verboso: bool) -> None:
        manipulador = logging.StreamHandler(sys.stderr)
        manipulador.addFilter(FiltroArquivo())
        manipulador.setFormatter(logging.Formatter("%(levelname)s [%(arquivo)s] %(message)s"))
        logging.basicConfig(
            level=logging.INFO if verboso else logging.WARNING, handlers=[manipulador], force=True
        )
    
    
    def criar_parser() -> argparse.ArgumentParser:
        parser = argparse.ArgumentParser(prog="processador")
        parser.add_argument("-v", "--verboso", action="store_true")
        sub = parser.add_subparsers(dest="comando", required=True)
    
        gerar = sub.add_parser("gerar", help="gera arquivos CSV de exemplo")
        gerar.add_argument("pasta", type=Path)
        gerar.add_argument("--arquivos", type=int, default=6)
        gerar.add_argument("--linhas", type=int, default=50_000)
    
        proc = sub.add_parser("processar", help="processa todos os CSV de uma pasta")
        proc.add_argument("pasta", type=Path)
        proc.add_argument("--leituras", type=int, default=4, help="leituras simultâneas")
        proc.add_argument("--processos", type=int, default=2, help="processos de cálculo")
        proc.add_argument("--timeout", type=float, default=5.0, help="segundos por arquivo")
        return parser
    
    
    def main(argv: list[str] | None = None) -> int:
        args = criar_parser().parse_args(argv)
        configurar_logs(args.verboso)
        if args.comando == "gerar":
            gerar_dados(args.pasta, args.arquivos, args.linhas)
            print(f"{args.arquivos} arquivos gerados em {args.pasta}")
            return 0
        with criar_pool(args.processos) as pool:
            relatorio = asyncio.run(
                processar(args.pasta, pool=pool, leituras=args.leituras, timeout=args.timeout)
            )
        imprimir(relatorio)
        return 1 if relatorio.falhas else 0
    
    projetos/processador/src/processador/__main__.py
    from processador.cli import main
    
    if __name__ == "__main__":
        raise SystemExit(main())
    

    Os testes

    O teste do cálculo é puro. O de pipeline cobre o que importa em um sistema concorrente: o arquivo corrompido não derruba os outros, as esperas crescem como combinado (0,1 e depois 0,2), o timeout vira uma Falha e o limite de leituras simultâneas é respeitado de verdade (medido com um contador protegido por trava):

    projetos/processador/tests/test_calculo.py
    import pytest
    
    from processador.calculo import agregar
    
    
    def test_agrega_por_regiao() -> None:
        texto = "regiao,valor_centavos\nsul,100\nnorte,50\nsul,25\n"
        assert agregar(texto) == (3, 175, {"sul": 125, "norte": 50})
    
    
    def test_cabecalho_errado_e_recusado() -> None:
        with pytest.raises(ValueError, match="cabeçalho"):
            agregar("a,b\n1,2\n")
    
    
    def test_linha_invalida_informa_o_numero_da_linha() -> None:
        with pytest.raises(ValueError, match="linha 3"):
            agregar("regiao,valor_centavos\nsul,100\nsul,abc\n")
    
    projetos/processador/tests/test_pipeline.py
    import asyncio
    import threading
    import time
    from pathlib import Path
    
    from processador.cli import gerar_dados
    from processador.pipeline import criar_pool, processar
    
    
    def rodar(pasta: Path, **opcoes):  # type: ignore[no-untyped-def]
        async def principal():  # type: ignore[no-untyped-def]
            with criar_pool(2) as pool:
                return await processar(pasta, pool=pool, **opcoes)
    
        return asyncio.run(principal())
    
    
    def test_processa_todos_e_soma(tmp_path: Path) -> None:
        gerar_dados(tmp_path, arquivos=3, linhas=100)
        relatorio = rodar(tmp_path)
        assert len(relatorio.resultados) == 3
        assert relatorio.falhas == []
        assert sum(r.linhas for r in relatorio.resultados) == 300
    
    
    def test_arquivo_corrompido_nao_derruba_os_outros(tmp_path: Path) -> None:
        gerar_dados(tmp_path, arquivos=2, linhas=10)
        (tmp_path / "quebrado.csv").write_text("regiao,valor_centavos\nsul,abc\n", encoding="utf-8")
        relatorio = rodar(tmp_path)
        assert len(relatorio.resultados) == 2
        assert [f.arquivo for f in relatorio.falhas] == ["quebrado.csv"]
        assert "linha 2" in relatorio.falhas[0].motivo
    
    
    def test_retentativas_com_espera_crescente(tmp_path: Path) -> None:
        gerar_dados(tmp_path, arquivos=1, linhas=5)
        esperas: list[float] = []
        falhas_restantes = {"n": 2}
    
        async def dormir_falso(segundos: float) -> None:
            esperas.append(segundos)
    
        def ler_instavel(caminho: Path) -> str:
            if falhas_restantes["n"] > 0:
                falhas_restantes["n"] -= 1
                raise OSError("disco ocupado")
            return caminho.read_text(encoding="utf-8")
    
        relatorio = rodar(tmp_path, dormir=dormir_falso, ler=ler_instavel)
        assert len(relatorio.resultados) == 1
        assert esperas == [0.1, 0.2]
    
    
    def test_timeout_vira_falha(tmp_path: Path) -> None:
        gerar_dados(tmp_path, arquivos=1, linhas=5)
    
        def ler_lento(caminho: Path) -> str:
            time.sleep(0.3)
            return caminho.read_text(encoding="utf-8")
    
        relatorio = rodar(tmp_path, timeout=0.05, ler=ler_lento)
        assert relatorio.resultados == []
        assert relatorio.falhas[0].motivo == "TimeoutError"
    
    
    def test_limite_de_leituras_simultaneas(tmp_path: Path) -> None:
        gerar_dados(tmp_path, arquivos=6, linhas=5)
        trava = threading.Lock()
        ativas = {"agora": 0, "maximo": 0}
    
        def ler_medindo(caminho: Path) -> str:
            with trava:
                ativas["agora"] += 1
                ativas["maximo"] = max(ativas["maximo"], ativas["agora"])
            time.sleep(0.05)
            with trava:
                ativas["agora"] -= 1
            return caminho.read_text(encoding="utf-8")
    
        rodar(tmp_path, leituras=2, ler=ler_medindo)
        assert ativas["maximo"] <= 2
    
    projetos/processador/tests/test_cli.py
    from pathlib import Path
    
    import pytest
    
    from processador.cli import main
    
    
    def test_gerar_e_processar(tmp_path: Path, capsys: pytest.CaptureFixture[str]) -> None:
        pasta = str(tmp_path / "dados")
        assert main(["gerar", pasta, "--arquivos", "3", "--linhas", "50"]) == 0
        assert main(["processar", pasta, "--processos", "2"]) == 0
        saida = capsys.readouterr().out
        assert "arquivos processados: 3  falhas: 0" in saida
        assert "total" in saida
    
    
    def test_codigo_de_saida_1_quando_ha_falha(tmp_path: Path) -> None:
        (tmp_path / "ruim.csv").write_text("x,y\n", encoding="utf-8")
        assert main(["processar", str(tmp_path), "--processos", "1"]) == 1
    

    Rodar

    Terminal
    cd projetos/processador
    uv sync
    uv run pytest
    uv run mypy
    uv run ruff check .
    
    Saída
    ..........                                                               [100%]
    10 passed
    Success: no issues found in 8 source files
    All checks passed!
    

    Gerando seis arquivos de cinquenta mil linhas e processando:

    Terminal
    uv run processador gerar dados --arquivos 6 --linhas 50000
    uv run processador processar dados
    
    Saída
    6 arquivos gerados em dados
    arquivos processados: 6  falhas: 0
      centro-oeste R$ 15.060.300,00
      nordeste     R$ 15.060.900,00
      norte        R$ 15.058.500,00
      sudeste      R$ 15.059.700,00
      sul          R$ 15.059.100,00
      total        R$ 75.298.500,00
    

    Com um arquivo corrompido na pasta, o programa conclui o resto, lista a falha e sai com o código 1:

    Resultado com um arquivo corrompido
    arquivos processados: 6  falhas: 1
      (as seis regiões, somando só os arquivos válidos)
      FALHA quebrado.csv: linha 2 inválida: ['sul', 'abc']
    

    Com -v, cada linha de log traz o nome do arquivo graças à ContextVar. A ordem das linhas muda de uma execução para outra, porque os arquivos terminam em tempos diferentes, e é isso mesmo:

    Logs com -v
    INFO [vendas_00.csv] processado: 50000 linhas
    INFO [vendas_01.csv] processado: 50000 linhas
    INFO [vendas_02.csv] processado: 50000 linhas
    

    Desafios

    1. Medir de verdade. Use o cProfile (capítulo 48) e compare --processos 1, 2 e 4 com arquivos maiores. Onde o ganho para de crescer, e por quê? (Dica: o custo de enviar o texto para outro processo e de iniciar cada um.)
    2. Streaming. Hoje cada arquivo é lido inteiro na memória. Reescreva agregar para processar por linhas, sem carregar o arquivo todo, e confira com tracemalloc a diferença de pico.
    3. Fila limitada. Troque o gather por uma fila com produtores e consumidores, de modo que a memória usada não cresça com o número de arquivos.
    4. Cancelamento. Trate o Ctrl+C: cancele as tarefas pendentes, espere as que estão em andamento e imprima um relatório parcial.
    5. Empacotar e publicar. Aplique o capítulo 49 e publique o pacote no TestPyPI.

    Por que este programa funciona no Windows e no macOS

    Esses sistemas criam processos com spawn, que reimporta o módulo principal. O __main__.py protege a execução com if __name__ == "__main__", e a função agregar está no nível do módulo. Sem os dois cuidados, cada processo filho tentaria rodar o programa inteiro de novo.

    O que a falta de estado compartilhado ganha

    Não há Lock neste programa. Cada tarefa produz um valor e devolve, e só o final os junta, em uma única thread. Quando você consegue desenhar a concorrência assim, evitando compartilhar memória, os bugs de corrida deixam de existir, em vez de serem tratados.