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.
| Requisito | Decisão | Capítulo |
|---|---|---|
| Ler muitos arquivos sem bloquear | asyncio com to_thread | 47 |
| Limitar leituras simultâneas | Semaphore | 47 |
| Calcular em paralelo de verdade | ProcessPoolExecutor | 46 |
| Tolerar falha de leitura | Retentativa com espera crescente | 51 |
| Não esperar para sempre | asyncio.timeout | 47 |
| Uma falha não derruba as outras | Resultado ou Falha, nunca exceção solta | 23 |
| Rastrear cada arquivo nos logs | ContextVar e filtro de log | 51 |
| Código testável | Injeção da espera e da leitura | 51 |
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.
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:
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:
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:
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:
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:
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:
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
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):
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")
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
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
cd projetos/processador
uv sync
uv run pytest
uv run mypy
uv run ruff check .
.......... [100%]
10 passed
Success: no issues found in 8 source files
All checks passed!
Gerando seis arquivos de cinquenta mil linhas e processando:
uv run processador gerar dados --arquivos 6 --linhas 50000
uv run processador processar dados
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:
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:
INFO [vendas_00.csv] processado: 50000 linhas
INFO [vendas_01.csv] processado: 50000 linhas
INFO [vendas_02.csv] processado: 50000 linhas
Desafios
- Medir de verdade. Use o
cProfile(capítulo 48) e compare--processos 1,2e4com 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.) - Streaming. Hoje cada arquivo é lido inteiro na memória. Reescreva
agregarpara processar por linhas, sem carregar o arquivo todo, e confira comtracemalloca diferença de pico. - Fila limitada. Troque o
gatherpor uma fila com produtores e consumidores, de modo que a memória usada não cresça com o número de arquivos. - Cancelamento. Trate o
Ctrl+C: cancele as tarefas pendentes, espere as que estão em andamento e imprima um relatório parcial. - 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__.pyprotege a execução comif __name__ == "__main__", e a funçãoagregarestá 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á
Lockneste 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.