Introdução à Abstração de Concorrência
A biblioteca padrão do Python oferece módulos como threading e multiprocessing para lidar com concorrência em nível baixo. A partir da versão 3.2, foi introduzido o módulo concurrent.futures, que encapsula essas implementações através de abstrações de alto nível. Este módulo fornece classes como ThreadPoolExecutor e ProcessPoolExecutor, simplificando significativamente o gerenciamento de pools de threads e processos.
O núcleo deste módulo repousa sobre dois pilares principais: Executores (Executor) e Futuros (Future).
A Classe Executor
Executor atua como uma classe base abstrata, não destinada a instâncias diretas, mas define a interface comum para executores concretos. Suas implementações específicas são responsáveis por alocar recursos e gerenciar o ciclo de vida das tarefas enviadas.
Abaixo, apresentamos um exemplo prático utilizando submit(). Diferentemente do exemplo clássico, utilizamos variáveis com nomes mais descritivos e alteramos a lógica interna da função simulada:
from concurrent.futures import ThreadPoolExecutor
import time
def simular_calculo_pesado(texto, atraso):
"""Função que representa uma tarefa demorada."""
time.sleep(atraso)
return f"Resultado: {texto}"
# Inicialização do pool com capacidade máxima de 2 workers
executor_thread = ThreadPoolExecutor(max_workers=2)
# Submissão de duas tarefas independentes
futuro_a = executor_thread.submit(simular_calculo_pesado, "PrimeiroTarefa", 2)
futuro_b = executor_thread.submit(simular_calculo_pesado, "SegundaTarefa", 3)
# Verifica se a tarefa já concluiu imediatamente após submissão
print(f"Tarefa A finalizada: {futuro_a.done()}")
# Aguarda tempo equivalente ao cálculo para verificar novamente
time.sleep(2.5)
print(f"Tarefa B finalizada: {futuro_b.done()}")
# Recuperação dos resultados garantindo que a execução terminou
print(futuro_a.result())
print(futuro_b.result())
executor_thread.shutdown(wait=True)
Para cenários de múltiplos núcleos de CPU onde o gargalo é computacional, basta substituir ThreadPoolExecutor por ProcessPoolExecutor.
Método Map para Processamento Ordenado
O método map() oferece uma maneira funcional de aplicar uma função a cada elemento de um ou mais iteráveis. Diferente do submit(), este método garante que os resultados sejam retornados na mesma ordem dos inputs fornecidos.
# coding: utf-8
from concurrent.futures import ThreadPoolExecutor
import requests
# Lista de endpoints para teste de requisição paralela
endereços_teste = [
'https://www.python.org',
'https://www.google.com',
'https://www.github.com'
]
def buscar_conteudo(url, tempo_limite=5):
try:
resp = requests.get(url, timeout=tempo_limite)
return {'url': url, 'tamanho': len(resp.content)}
except Exception:
return {'url': url, 'erro': True}
with ThreadPoolExecutor(max_workers=3) as pool_exec:
# Mapeamento das funções para a lista de URLs
resultados_ordenados = pool_exec.map(buscar_conteudo, endereços_teste)
for dado in resultados_ordenados:
print(f"{dado['url']}: {'Erro' if dado.get('erro') else str(dado['tamanho'])} bytes")
Objetos Future e Estado das Tarefas
O conceito de Future é fundamental para programação assíncrona. Ele repreesnta um resultado que pode não estar disponível imediatamente. Enquanto o sistema aguarda a conclusão de operações I/O, o processador não fica ocioso, podendo gerenciar outros estados.
Quando usamos submit(), recebemos uma instância Future que permite inspecionar o estado atual da operação antes mesmo de obter o valor final.
from concurrent.futures import ThreadPoolExecutor, as_completed
import requests
objetivos_web = ['https://bing.com', 'https://yahoo.com', 'https://example.com']
def requisicao_web(url, timeout=5):
return requests.get(url, timeout=timeout)
with ThreadPoolExecutor(max_workers=3) as executor:
# Criação de todos os futuros
futurus_processo = [executor.submit(requisicao_web, site) for site in objetivos_web]
# Monitoramento inicial do estado de execução
for f in futurus_processo:
if f.running():
print(f"Iniciando processamento para: {f}")
# Iteração assim que cada tarefa finalizar, independente da ordem
for f in as_completed(futurus_processo):
try:
resposta = f.result()
print(f"Sucesso - URL: {resposta.url}, Tamanho: {len(resposta.content)} bytes")
except Exception as excecao:
print(f"Falha na tarefa: {excecao}")
O uso de as_completed() retorna os objetos futuros conforme eles terminam, permitindo processar os dados mais rápidos sem esperar pelos lentos, diferentemente do comportamento ordenado do método map.
Sincronização com Espera Controlada (wait)
A função wait() bloqueia a execução até que uma condição específica seja atendida por um conjunto de Future objects. Ela retorna uma tupla contendo dois conjuntos: tarefas concluídas (done_set) e tarefas pendentes (pending_set).
Você pode ajustar o comportamento de retorno usando o parâmetro return_when:
ALL_COMPLETED: Espera até que todas as tarefas estejam feitas (padrão).FIRST_COMPLETED: Retorna assim que a primeira tarefa terminar.FIRST_EXCEPTION: Retorna se alguma tarefa levantar uma exceção.
Também é possível estabelecer um limite de tempo máximo via parâmetro timeout.
from concurrent.futures import ThreadPoolExecutor, wait, FIRST_COMPLETED
lista_sites = ['https://site1.com', 'https://site2.com']
def carregar_site(url):
import time
time.sleep(2)
return url + "_carregado"
com executor := ThreadPoolExecutor(max_workers=2) as e:
tarefas = [e.submit(carregar_site, s) for s in lista_sites]
# Aguarda até que pelo menos uma tarefa termine (FIRST_COMPLETED) ou exceda o tempo
concluidos, pendentes = wait(tarefas, timeout=5, return_when=FIRST_COMPLETED)
print(f"Tarefas prontas: {len(concluidos)}")
print(f"Tarefas em espera: {len(pendentes)}")