Implementação de Tarefas Agendadas com Redis ZSET e List

Estratégia de Implementação

A solução baseia-se na combinação de duas estruturas de dados do Redis: ZSET para tarefas com execução futura e LIST para tarefas imediatas. O campo score do ZSET armazena timestamps em milissegundos, representando o instante exato de execução. Um processo periódico compara o tempo atual com os scores armazenados — toda entrada cujo score seja menor ou igual ao tempo corrente é migrada para uma lista de execução imediata.

Para tarefas que devem ser processadas assim que enfileiradas (ex.: eventos de alta prioridade), o dado é inserido diretamente em uma LIST, evitando qualquer avaliação temporal no lado do Redis.

Justificativa da Escolha entre LIST e ZSET

  • LIST: Estrutura de dupla ligação com inserção/remoção em O(1) nas extremidades. Ideal para filas FIFO/LIFO de alto throughput, especialmente quando a ordem de chegada é suficiente e não há necessidade de ordenação por tempo.
  • ZSET: Permite ordenação por score com busca eficeinte por intervalos (ZRANGEBYSCORE). Adequado para agendamento preciso, mas com custo maior em operações de remoção individual comparado à LIST.

Estratégia de Gerenciamento de Dados

Para evitar sobrecarga na memória do Redis e garantir estabilidade, adota-se uma política de *time-window slicing*: apenas tarefas com tempo de execução dentro dos próximos 5 minutos são mantidas no ZSET. Tarefas mais distantes são persistidas no banco relacional (ex.: MySQL), reduzindo a superfície de risco de bloqueio ou lentidão causada por grandes conjuntos ordenados.

Sincronização entre Banco de Dados e Redis

O sistema implementa dois fluxos de sincronização:

  1. Migração de tarefas iminentes: Um job agendado a cada minuto varre todos os ZSETs (usando SCAN com padrão "future:*"), extrai itens com score ≤ now() e os transfere atomicamente para suas respectivas listas de consumo via pipeline.
  2. Carga inicial de tarefas futuras: A cada 5 minutos, o serviço consulta o banco por tarefas com execute_time < now() + 5min e as insere no ZSET correspondente — após limpeza prévia dos dados em cache para evitar duplicações.

Modelos de Domínio Reestruturados

Os POJOs foram reescritos para maior clareza e alinhamento com boas práticas de serialização segura:

package com.example.scheduler.model;

import java.io.Serializable;
import java.time.Instant;

public class ScheduledJob implements Serializable {
    private static final long serialVersionUID = 42L;

    private Long jobId;
    private Integer category;
    private Integer urgencyLevel;
    private Instant scheduledAt;
    private byte[] payload;

    // getters e setters omitidos para concisão
}

package com.example.scheduler.model;

import com.baomidou.mybatisplus.annotation.IdType;
import com.baomidou.mybatisplus.annotation.TableId;
import com.baomidou.mybatisplus.annotation.TableName;
import java.io.Serializable;
import java.time.LocalDateTime;

@TableName("job_schedule")
public class JobScheduleRecord implements Serializable {
    private static final long serialVersionUID = 43L;

    @TableId(type = IdType.ASSIGN_ID)
    private Long jobId;
    private LocalDateTime scheduledAt;
    private byte[] payload;
    private Integer category;
    private Integer urgencyLevel;
}

package com.example.scheduler.model;

import com.baomidou.mybatisplus.annotation.*;
import java.io.Serializable;
import java.time.LocalDateTime;

@TableName("job_execution_log")
public class JobExecutionLog implements Serializable {
    private static final long serialVersionUID = 44L;

    @TableId(type = IdType.ASSIGN_ID)
    private Long jobId;
    private LocalDateTime scheduledAt;
    private byte[] payload;
    private Integer category;
    private Integer urgencyLevel;
    @Version
    private Integer revision;
    private Integer status; // 0=PENDENTE, 1=CONCLUÍDO, 2=CANCELADO
}

Serviço de Agendamento com Sincronização Distribuída

O serviço emprega bloqueio distribuído via SETNX para garantir que apenas uma instância execute a migração de tarefas do ZSET para a LIST simultaneamente. Também utiliza otimismo com versionamento para evitar execuções duplicadas de tarefas críticas.

@Service
public class JobSchedulerService {

    @Autowired private CacheClient redis;
    @Autowired private JobScheduleMapper scheduleMapper;
    @Autowired private JobLogMapper logMapper;

    private static final String FUTURE_PREFIX = "future:";
    private static final String QUEUE_PREFIX = "queue:";

    public void enqueue(ScheduledJob job) {
        scheduleMapper.insert(toScheduleRecord(job));

        final long now = System.currentTimeMillis();
        final long deadline = now + TimeUnit.MINUTES.toMillis(5);

        final String key = buildKey(job.getCategory(), job.getUrgencyLevel());

        if (job.getScheduledAt().toEpochMilli() <= now) {
            redis.lpush(QUEUE_PREFIX + key, serialize(job));
        } else if (job.getScheduledAt().toEpochMilli() <= deadline) {
            redis.zadd(FUTURE_PREFIX + key, job.getScheduledAt().toEpochMilli(), serialize(job));
        }
    }

    @Scheduled(cron = "0 */1 * * * ?")
    public void migrateImminentJobs() {
        final String lockKey = "lock:migrate";
        final String token = redis.tryAcquireLock(lockKey, 30_000);
        if (token == null) return;

        try {
            redis.scan(FUTURE_PREFIX + "*").forEach(futureKey -> {
                final String queueKey = QUEUE_PREFIX + futureKey.substring(FUTURE_PREFIX.length());
                final Set<String> dueJobs = redis.zrangeByScore(futureKey, 0, System.currentTimeMillis());
                if (!dueJobs.isEmpty()) {
                    redis.pipeline()
                         .zrem(futureKey, dueJobs.toArray(new String[0]))
                         .lpush(queueKey, dueJobs.toArray(new String[0]))
                         .exec();
                }
            });
        } finally {
            redis.releaseLock(lockKey, token);
        }
    }

    @Scheduled(cron = "0 */5 * * * ?")
    public void syncUpcomingJobs() {
        final LocalDateTime cutoff = LocalDateTime.now().plusMinutes(5);
        final List<JobScheduleRecord> upcoming = scheduleMapper.selectList(
            Wrappers.<JobScheduleRecord>lambdaQuery()
                .le(JobScheduleRecord::getScheduledAt, cutoff)
        );

        upcoming.forEach(record -> {
            final ScheduledJob job = fromRecord(record);
            final String key = buildKey(job.getCategory(), job.getUrgencyLevel());
            redis.zadd(FUTURE_PREFIX + key,
                       job.getScheduledAt().toEpochMilli(),
                       serialize(job));
        });
    }

    private String buildKey(Integer category, Integer level) {
        return category + ":" + level;
    }

    private String serialize(ScheduledJob job) {
        return JacksonUtil.toJson(job);
    }

    private ScheduledJob fromRecord(JobScheduleRecord record) {
        ScheduledJob job = new ScheduledJob();
        job.setJobId(record.getJobId());
        job.setCategory(record.getCategory());
        job.setUrgencyLevel(record.getUrgencyLevel());
        job.setScheduledAt(record.getScheduledAt().atZone(ZoneId.systemDefault()).toInstant());
        job.setPayload(record.getPayload());
        return job;
    }

    private JobScheduleRecord toScheduleRecord(ScheduledJob job) {
        JobScheduleRecord record = new JobScheduleRecord();
        record.setJobId(job.getJobId());
        record.setCategory(job.getCategory());
        record.setUrgencyLevel(job.getUrgencyLevel());
        record.setScheduledAt(LocalDateTime.ofInstant(job.getScheduledAt(), ZoneId.systemDefault()));
        record.setPayload(job.getPayload());
        return record;
    }
}

Tags: Redis Zset List scheduling distributed-lock

Publicado em 9-22 14:33