concurrent.futures
Nesta aula, você aprenderá a usar o módulo concurrent.futures do Python para executar tarefas em paralelo com ThreadPoolExecutor e ProcessPoolExecutor, entender o conceito de Futures e aplicar padrões comuns de concorrência para melhorar o desempenho de seus programas.
O módulo concurrent.futures é uma das ferramentas mais poderosas e fáceis de usar do Python para programação concorrente. Ele fornece uma API de alto nível para executar chamadas de funções de forma assíncrona, usando threads ou processos, sem que você precise gerenciar diretamente a criação, sincronização e comunicação entre eles. Isso torna o código mais limpo, mais seguro e mais portátil.
Nesta aula, vamos explorar os dois principais executores: ThreadPoolExecutor e ProcessPoolExecutor. Também vamos entender o conceito de Futures, que são objetos que representam o resultado futuro de uma operação assíncrona, e veremos padrões de uso comuns, como executar tarefas em lote, processar resultados à medida que ficam prontos e lidar com exceções.
ThreadPoolExecutor
O ThreadPoolExecutor utiliza um pool de threads para executar chamadas de função. Ele é ideal para tarefas que são limitadas por I/O (input/output), como operações de rede, leitura/escrita de arquivos, ou chamadas a APIs externas, onde o tempo de espera é significativo. Como as threads compartilham o mesmo espaço de memória, a comunicação é simples, mas é preciso ter cuidado com a concorrência em dados compartilhados, pois o GIL (Global Interpreter Lock) do Python não protege contra corridas em operações que não sejam atômicas.
Para usar o ThreadPoolExecutor, você cria uma instância, especificando o número máximo de threads (padrão é o número de processadores da máquina, mas para I/O, pode ser maior). Em seguida, você submete funções usando o método submit(), que retorna um Future, ou usa o método map() para aplicar uma função a uma lista de argumentos de forma concorrente. O contexto with garante que o executor seja encerrado corretamente.
import concurrent.futures
import time
def tarefa(nome, segundos):
print(f"Iniciando {nome}...")
time.sleep(segundos)
print(f"Terminando {nome}...")
return f"{nome} levou {segundos}s"
# Criando um ThreadPoolExecutor com 3 threads
with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
# Submetendo tarefas individualmente
futures = [executor.submit(tarefa, f"Tarefa-{i}", i) for i in range(1, 5)]
# Coletando resultados na ordem de conclusão
for future in concurrent.futures.as_completed(futures):
print(future.result())
# Saída (ordem pode variar):
# Iniciando Tarefa-1...
# Iniciando Tarefa-2...
# Iniciando Tarefa-3...
# Terminando Tarefa-1...
# Tarefa-1 levou 1s
# Terminando Tarefa-2...
# Tarefa-2 levou 2s
# Terminando Tarefa-3...
# Tarefa-3 levou 3s
# Iniciando Tarefa-4...
# Terminando Tarefa-4...
# Tarefa-4 levou 4sNo exemplo acima, usamos as_completed() para processar os resultados assim que cada tarefa termina, independentemente da ordem de submissão. Isso é eficiente quando as tarefas têm durações variáveis. Se você quiser os resultados na ordem de submissão, pode usar executor.map(), que retorna um iterador que produz os resultados na mesma ordem dos argumentos fornecidos.
ProcessPoolExecutor
O ProcessPoolExecutor usa processos separados, cada um com seu próprio interpretador Python e espaço de memória. Isso contorna o GIL e permite que tarefas com uso intensivo de CPU (processamento numérico, cálculos complexos) sejam executadas verdadeiramente em paralelo, aproveitando múltiplos núcleos. No entanto, a comunicação entre processos é mais cara, pois exige serialização dos dados (via pickle), e o custo de criar processos é maior.
Para usar o ProcessPoolExecutor, a sintaxe é quase idêntica à do ThreadPoolExecutor. Uma diferença importante é que as funções e os argumentos devem ser serializáveis (picklable). Além disso, em sistemas Windows, é necessário proteger o código principal com if __name__ == '__main__': para evitar problemas com a criação de processos.
import concurrent.futures
import math
def calcular_fatorial(n):
return math.factorial(n)
# Criando um ProcessPoolExecutor com 4 processos
if __name__ == '__main__':
numeros = [100000, 200000, 300000, 400000]
with concurrent.futures.ProcessPoolExecutor(max_workers=4) as executor:
resultados = list(executor.map(calcular_fatorial, numeros))
print(resultados) # Exibe os fatoriais calculadosNo exemplo, usamos executor.map() para aplicar a função calcular_fatorial a cada número. Como o cálculo de fatorial é intensivo em CPU, o uso de processos pode trazer ganho significativo de desempenho em máquinas com múltiplos núcleos. É importante notar que o ProcessPoolExecutor é mais adequado para tarefas que não exigem muita comunicação entre processos, pois cada chamada é independente.
Futures
Um Future é um objeto que representa o resultado de uma operação assíncrona. Ele é retornado pelo método submit() e oferece métodos para verificar se a tarefa foi concluída, aguardar a conclusão e obter o resultado ou a exceção. O módulo concurrent.futures fornece a classe Future com os seguintes métodos principais:
result(timeout=None): retorna o resultado da chamada. Se a chamada ainda não estiver concluída, aguarda até o timeout (ou indefinidamente se None). Se a chamada levantou uma exceção, ela é relançada aqui.exception(timeout=None): retorna a exceção levantada pela chamada, ou None se não houve exceção.done(): retornaTruese a chamada foi concluída (com sucesso ou com exceção).add_done_callback(fn): adiciona uma função de retorno que será chamada quando o Future for concluído.
Além disso, as funções de conveniência concurrent.futures.wait() e concurrent.futures.as_completed() permitem aguardar vários Futures de formas diferentes.
import concurrent.futures
import time
def tarefa_demorada():
time.sleep(2)
return 42
with concurrent.futures.ThreadPoolExecutor() as executor:
future = executor.submit(tarefa_demorada)
print(f"A tarefa está concluída? {future.done()}") # False
# Aguardando com timeout
try:
resultado = future.result(timeout=1)
except concurrent.futures.TimeoutError:
print("A tarefa ainda não terminou")
# Aguardando indefinidamente
resultado = future.result()
print(f"Resultado final: {resultado}")
# Usando callback
def callback(fut):
print(f"Callback chamado com {fut.result()}")
future2 = executor.submit(tarefa_demorada)
future2.add_done_callback(callback)Os Futures são a base para a programação assíncrona em Python. Eles permitem que você escreva código que não bloqueia a execução enquanto espera por resultados, e podem ser combinados para criar fluxos de trabalho complexos.
Padrões
Existem vários padrões comuns de uso do concurrent.futures que você encontrará em aplicações reais. Aqui estão alguns dos mais úteis:
Padrão 1: Mapa paralelo
Use executor.map() quando você tem uma função e uma lista de argumentos, e deseja aplicar a função a todos os argumentos em paralelo, coletando os resultados na mesma ordem. Isso é semelhante à função map() embutida, mas executa as chamadas de forma concorrente.
import concurrent.futures
import urllib.request
def baixar_url(url):
with urllib.request.urlopen(url) as response:
return len(response.read())
urls = [
'https://www.python.org',
'https://www.wikipedia.org',
'https://www.github.com',
]
with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
tamanhos = list(executor.map(baixar_url, urls))
print(tamanhos) # Tamanho de cada página em bytesPadrão 2: Processar resultados à medida que chegam
Quando você precisa processar resultados na ordem em que são concluídos, use as_completed(). Isso é útil quando as tarefas têm durações variáveis e você quer começar a processar o resultado mais cedo possível.
import concurrent.futures
import time
def tarefa(n):
time.sleep(n)
return n
with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
futures = [executor.submit(tarefa, i) for i in [3, 1, 2]]
for future in concurrent.futures.as_completed(futures):
print(future.result()) # Ordem de conclusão: 1, 2, 3Padrão 3: Tratamento de exceções
Quando uma tarefa levanta uma exceção, o Future captura essa exceção e a relança quando você chama result() ou exception(). Você pode tratar exceções individualmente ou agrupá-las.
import concurrent.futures
def dividir(a, b):
return a / b
with concurrent.futures.ThreadPoolExecutor() as executor:
futures = [executor.submit(dividir, 10, i) for i in [2, 0, 5]]
for future in futures:
try:
resultado = future.result()
print(f"Resultado: {resultado}")
except ZeroDivisionError:
print("Divisão por zero!")
except Exception as e:
print(f"Erro: {e}")Padrão 4: Combinação de múltiplos Futures
Você pode aguardar vários Futures com wait() e depois coletar os resultados. Isso é útil quando você precisa que todas as tarefas terminem antes de continuar.
import concurrent.futures
import time
def tarefa(n):
time.sleep(n)
return n * 2
with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor:
futures = [executor.submit(tarefa, i) for i in [2, 3, 1]]
done, not_done = concurrent.futures.wait(futures, return_when=concurrent.futures.ALL_COMPLETED)
resultados = [f.result() for f in done]
print(resultados) # [4, 6, 2]Boas Práticas e Observações Finais
- Escolha o executor certo: Use
ThreadPoolExecutorpara tarefas limitadas por I/O eProcessPoolExecutorpara tarefas limitadas por CPU. Não use processos para tarefas que fazem muitas operações de I/O, pois o overhead de serialização pode dominar. - Limite o número de workers: Defina
max_workerscom base no tipo de tarefa. Para I/O, você pode ter mais threads do que núcleos, mas para CPU, use o número de núcleos (ou um pouco mais, se houver esperas). - Evite compartilhar estado: Em threads, evite modificar variáveis globais sem sincronização. Em processos, o estado não é compartilhado, então você deve passar dados via argumentos e retornos.
- Use contextos
with: Isso garante que o executor seja encerrado adequadamente, liberando recursos. - Trate exceções: Sempre chame
future.result()dentro de um try/except para capturar exceções que possam ocorrer na tarefa.
Referências
- Documentação oficial do módulo concurrent.futures
- Documentação do módulo threading
- Documentação do módulo multiprocessing
- Artigo sobre concorrência em Python no Real Python
- Concurrent futures em Python - GeeksforGeeks
- Python Module of the Week: concurrent.futures
Exercícios
- Exercício 1: Crie um programa que use
ThreadPoolExecutorpara baixar o conteúdo de 5 URLs (useurllib.request) e imprima o tamanho de cada conteúdo. Useas_completed()para imprimir os tamanhos na ordem em que as respostas chegarem.✓ Resposta:import concurrent.futures import urllib.request urls = [ 'https://www.python.org', 'https://www.wikipedia.org', 'https://www.github.com', 'https://www.stackoverflow.com', 'https://www.google.com', ] def baixar(url): with urllib.request.urlopen(url) as response: return url, len(response.read()) with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor: futures = [executor.submit(baixar, url) for url in urls] for future in concurrent.futures.as_completed(futures): url, tamanho = future.result() print(f"{url}: {tamanho} bytes") - Exercício 2: Escreva uma função que calcule o quadrado de um número. Use
ProcessPoolExecutorpara calcular os quadrados de uma lista de números de 1 a 10 e imprima os resultados em ordem crescente. Lembre-se de proteger o código principal comif __name__ == '__main__':.✓ Resposta:import concurrent.futures def quadrado(x): return x * x if __name__ == '__main__': numeros = list(range(1, 11)) with concurrent.futures.ProcessPoolExecutor(max_workers=4) as executor: resultados = list(executor.map(quadrado, numeros)) print(resultados) # [1, 4, 9, 16, 25, 36, 49, 64, 81, 100] - Exercício 3: Usando
ThreadPoolExecutor, crie um programa que simule o processamento de pedidos. Cada pedido leva um tempo aleatório (1-5 segundos) para ser processado. Usesubmit()ewait()para aguardar todos os pedidos e imprima o tempo total gasto.✓ Resposta:import concurrent.futures import time import random def processar_pedido(pedido): tempo = random.randint(1, 5) time.sleep(tempo) return f"Pedido {pedido} processado em {tempo}s" if __name__ == '__main__': inicio = time.time() pedidos = [1, 2, 3, 4, 5] with concurrent.futures.ThreadPoolExecutor(max_workers=3) as executor: futures = [executor.submit(processar_pedido, p) for p in pedidos] done, _ = concurrent.futures.wait(futures, return_when=concurrent.futures.ALL_COMPLETED) for f in done: print(f.result()) fim = time.time() print(f"Tempo total: {fim - inicio:.2f}s") - Exercício 4: Escreva um programa que use
ProcessPoolExecutorpara calcular a soma de uma lista de números, dividindo a lista em partes. Cada processo soma uma parte e o resultado final é a soma das partes. Use uma lista de 1 milhão de números (crie comrange(1, 1000001)).✓ Resposta:import concurrent.futures def soma_parte(lista): return sum(lista) if __name__ == '__main__': numeros = list(range(1, 1000001)) # Dividir em 4 partes tamanho = len(numeros) // 4 partes = [numeros[i*tamanho:(i+1)*tamanho] for i in range(4)] # Ajustar a última parte para incluir o restante partes[-1] += numeros[4*tamanho:] with concurrent.futures.ProcessPoolExecutor(max_workers=4) as executor: resultados = list(executor.map(soma_parte, partes)) total = sum(resultados) print(f"Soma total: {total}") - Exercício 5: Crie um programa que use
ThreadPoolExecutorpara fazer requisições HTTP GET a uma API pública (por exemplo,https://jsonplaceholder.typicode.com/posts/1até/5) e imprima o título de cada post. Useadd_done_callback()para processar cada resposta assim que ela chegar.✓ Resposta:import concurrent.futures import urllib.request import json def buscar_post(id): url = f"https://jsonplaceholder.typicode.com/posts/{id}" with urllib.request.urlopen(url) as response: data = json.loads(response.read().decode()) return data['title'] def imprimir_titulo(future): print(f"Título: {future.result()}") with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor: futures = [executor.submit(buscar_post, i) for i in range(1, 6)] for f in futures: f.add_done_callback(imprimir_titulo)