Antes da existência de thread pools, a criação de threads em Java se limitava a duas abordagens: estender Thread ou implementar Runnable. Embora funcionais, essas abordagens apresentavam limitações significativas:
- Overhead considerável na criação e destruição de threads, impactando diretamente a performance
- Risco de exaustão de memória com criação descontrolada
- Impossibilidade de reutilização imediata de threads para novas tarefas
A solução veio com o conceito de pool de threads: um repositório gerenciado de threads reutilizáveis. Os benefícios incluem economia de recursos, resposta mais ágil a novas demandas e controle cenrtalizado sobre o ciclo de vida das threads.
Arquitetura Interna
O construtor ThreadPoolExecutor revela os componentes essenciais:
public ThreadPoolExecutor(
int tamanhoBase,
int tamanhoMaximo,
long tempoInatividade,
TimeUnit unidadeTemporal,
BlockingQueue<Runnable> filaTarefas,
ThreadFactory fabricaThreads,
RejectedExecutionHandler tratadorRejeicoes)
Os parâmetros definem: quantidade mínima e máxima de threads, tempo de sobrevivência de threads excedentes, fila de espera para tarefas, fabricação personalizada de threads e estratégia para tarefas rejeitadas.
O fluxo operacional segue esta lógica: inicialização com threads base, enfileiramento quando o limite é atingido, expansão até o máximo se a fila saturar, reciclagem automática de threads ociosas, e aplicação de políticas de rejeição em cenários extremos.
Armadilhas Comuns
O utilitário Executors oferece métodos de conveniência que escondem riscos:
Fila ilimitada: newFixedThreadPool emprega LinkedBlockingQueue sem limite explícito, podendo consumir toda a memória disponível com acúmulo de tarefas.
// Implementação problemática - capacidade implícita de Integer.MAX_VALUE
public static ExecutorService newFixedThreadPool(int nThreads) {
return new ThreadPoolExecutor(nThreads, nThreads,
0L, TimeUnit.MILLISECONDS,
new LinkedBlockingQueue<Runnable>());
}
Explosão de threads: newCachedThreadPool permite crescimento ilimitado da quantidade de threads, com SynchronousQueue que não armazena elementos.
// Risco de criação massiva de threads
public static ExecutorService newCachedThreadPool() {
return new ThreadPoolExecutor(0, Integer.MAX_VALUE,
60L, TimeUnit.SECONDS,
new SynchronousQueue<Runnable>());
}
Perda de dados em memória: o problema mais crítico. Como as filas residem na memória RAM, qualquer interrupção do processo — reinicialização, queda de energia, kill do container — elimina tarefas pendentes irreversivelmente.
Estratégia de Persistência Garantida
Para eliminar a perda de dados, é necessário externalizar o estado das tarefas. A solução arquitetural envolve:
1. Escrita transacional: ao receber uma requisição, execute a operação crítica síncrona e registre uma entrada em banco de dados marcada como pendente, tudo dentro da mesma transação.
2. Processamento assíncrono orquestrado: um mecanismo de polling periodicamente recupera registros pendentes ordenados cronologicamente, submetendo-os ao pool de threads.
3. Atualização de estado: após processamento bem-sucedido, o registro transita para concluído. Em caso de falha, incrementa-se um contador de tentativas.
Esta abordagem exige idempotência na lógica de processamento — múltiplas execuções da mesma tarefa devem produzir resultado equivalente à execução única.
Implementação de Referência
@Component
public class ProcessadorTarefasPersistidas {
private final TarefaRepositorio repositorio;
private final ExecutorService executor;
private final int maximoTentativas = 3;
public ProcessadorTarefasPersistidas(TarefaRepositorio repo) {
this.repositorio = repo;
this.executor = new ThreadPoolExecutor(
4, 8, 60, TimeUnit.SECONDS,
new ArrayBlockingQueue<>(100),
new CustomThreadFactory(),
new ThreadPoolExecutor.CallerRunsPolicy()
);
}
@Scheduled(fixedRate = 30000)
public void despacharPendentes() {
Pageable paginacao = PageRequest.of(0, 50, Sort.by("criadoEm"));
List<Tarefa> candidatas = repositorio
.findByStatusAndTentativasLessThan(Status.PENDENTE, maximoTentativas, paginacao);
candidatas.forEach(tarefa -> executor.submit(() -> executarComSeguranca(tarefa)));
}
private void executarComSeguranca(Tarefa tarefa) {
try {
processadorNegocio.executar(tarefa.getPayload());
repositorio.atualizarStatus(tarefa.getId(), Status.CONCLUIDO);
} catch (Exception ex) {
repositorio.incrementarTentativas(tarefa.getId());
if (tarefa.getTentativas() + 1 >= maximoTentativas) {
repositorio.atualizarStatus(tarefa.getId(), Status.FALHA_CRITICA);
alertarMonitoramento(tarefa, ex);
}
}
}
}
O esquema de tabela suporta esta operação:
CREATE TABLE tarefas_assincronas (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
payload JSON NOT NULL,
status ENUM('PENDENTE', 'CONCLUIDO', 'FALHA_CRITICA') DEFAULT 'PENDENTE',
tentativas INT DEFAULT 0,
criado_em TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
atualizado_em TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
INDEX idx_status_tentativas (status, tentativas, criado_em)
);
Com esta estrutura, mesmo que o serviço seja interrompido abruptamente, as tarefas permanecem no banco de dados. Na próxima ativação do scheduler, elas serão reprocessadsa atuomaticamente, garantindo durabilidade sem comprometer a assincronicidade.