Preservação de Tarefas em Thread Pools Durante Falhas de Serviço

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.

Tags: ThreadPoolExecutor Java Concurrency Task Persistence Database Queue Idempotency

Publicado em 10-5 05:15