Gerenciamento de Execução Assíncrona com concurrent.futures em Python

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)}")

Tags: Python concurrent-futures Multi-threading multi-processing future-object

Publicado em 8-15 04:49