Implementando um Agendador de Tarefas em Python similar ao Crontab do Linux

Introdução

Muitas vezes precisamos executar tarefas em intervalos específicos ou com timeout controlado. O Linux oferece o Crontab para isso, mas podemos impelmentar nossa própria solução em Python para ter mais flexibilidade e controle sobre o ambiente de execução.

Neste artigo, vamos explorar como criar um sistema de agendamento de tarefas em Python que permite:

  • Executar processos com tempo máximo de execução
  • Reiniciar automaticamente após um intervalo
  • Gerenciar múltiplas tarefas em paralelo

Requisitos do Sistema

Vamos definir os requisitos básicos para nosso agendador:

  • Um programa que permanece em execução contínua
  • Capacidade de terminate o processo após X segundos
  • Reinício automático após um intervalo definido

Implementação Básica

A seguir, apresentamos uma implementação que atende aos requisitos fundamentais:

#!/usr/bin/env python
# -*- coding=utf-8 -*-

import sys
import time
import multiprocessing
import datetime
import os
from typing import Optional, Union


class AgendadorError(Exception):
    """Exceção personalizada para erros de agendamento."""
    
    def __init__(self, mensagem: Optional[str] = None, codigo: int = 4999):
        self.codigo = int(codigo)
        mensagem = mensagem or "O tempo agendado é menor que o tempo atual"
        Exception.__init__(self, self.codigo, mensagem)


class ConversorData:
    """Converte diferentes formatos de data para timestamp."""
    
    @staticmethod
    def para_timestamp(data: Union[datetime.datetime, str, float]) -> float:
        """Converte data para timestamp Unix."""
        if isinstance(data, datetime.datetime):
            return time.mktime(data.timetuple())
        
        elif isinstance(data, str):
            formato = "%Y-%m-%d %H:%M:%S" if " " in data else "%Y-%m-%d"
            dt_obj = datetime.datetime.strptime(data.strip(), formato)
            return time.mktime(dt_obj.timetuple())
        
        elif isinstance(data, float):
            return data
        
        raise ValueError("Formato de data inválido")


class Temporizador:
    """Gerencia contagem regressiva para timeout de tarefas."""
    
    @staticmethod
    def executar(segundos: int) -> bool:
        """Executa contagem regressiva."""
        contador = segundos
        while contador > 0:
            sys.stdout.write('\r')
            contador -= 1
            sys.stdout.write(f"{contador}")
            sys.stdout.flush()
            time.sleep(1)
        return True


class Tarefas:
    """Demonstração de tarefas a serem executadas."""
    
    @staticmethod
    def iniciar():
        """Simula início de uma tarefa."""
        print("Tarefa iniciada")
        time.sleep(2)
        print("Tarefa concluída")
    
    @staticmethod
    def limpar():
        """Limpa recursos da tarefa."""
        print("Recursos limpos")


def processo_trabalho():
    """Função executada pelo processo filho."""
    print('Iniciando processamento')
    time.sleep(0.1)
    print('Processamento finalizado')
    Tarefas.iniciar()


def aguardar_horario(horario: str) -> bool:
    """Aguarda até o horário especificado."""
    timestamp_alvo = ConversorData.para_timestamp(horario)
    timestamp_atual = time.time()
    diferenca = int(timestamp_alvo - timestamp_atual)

    if diferenca < 0:
        raise AgendadorError()

    while diferenca > 0:
        diferenca -= 1
        time.sleep(1)
    return True


def executar_agendamento(
    horario: str = '',
    repeticoes: int = 1,
    timeout: int = 0,
    intervalo_reinicio: int = 0
) -> None:
    """
    Executa o agendamento de tarefas.
    
    Args:
        horario: Horário no formato "YYYY-MM-DD HH:MM:SS"
        repeticoes: Número de vezes para executar
        timeout: Tempo máximo de execução em segundos
        intervalo_reinicio: Intervalo para reinício em segundos
    """
    if horario:
        await_status = aguardar_horario(horario)
        if not await_status:
            return

    for _ in range(repeticoes):
        processo = multiprocessing.Process(target=processo_trabalho)
        processo.start()
        
        if timeout > 0:
            status = Temporizador.executar(timeout)
            if status:
                processo.terminate()
                Tarefas.limpar()
        
        if intervalo_reinicio > 0:
            time.sleep(intervalo_reinicio)


if __name__ == '__main__':
    executar_agendamento(repeticoes=1, timeout=50)

Versão com Suporte a Múltiplas Tarefas

Para cenários onde precisamos executar várias tarefas simultaneamente, podemos implementar um pool de threads:

#!/usr/bin/env python
# -*- coding=utf-8 -*-

import queue
import threading
import contextlib
import sys
import time
import multiprocessing
import datetime
import os
from typing import Optional, Union, Callable, Tuple, Any


class ErroAgendador(Exception):
    """Exceção para erros do agendador."""
    
    def __init__(self, msg: Optional[str] = None, codigo: int = 4999):
        self.codigo = int(codigo)
        msg = msg or "Erro desconhecido"
        Exception.__init__(self, self.codigo, msg)


class ConversorDataHora:
    """Utilitário para conversão de datas."""
    
    @staticmethod
    def converter_timestamp(data: Union[datetime.datetime, str, float]) -> float:
        if isinstance(data, datetime.datetime):
            return time.mktime(data.timetuple())
        
        elif isinstance(data, str):
            fmt = "%Y-%m-%d %H:%M:%S" if " " in data else "%Y-%m-%d"
            dt = datetime.datetime.strptime(data.strip(), fmt)
            return time.mktime(dt.timetuple())
        
        elif isinstance(data, float):
            return data
        
        raise ValueError("Formato inválido")


class ContadorRegressivo:
    """Implementa contagem regressiva."""
    
    @staticmethod
    def iniciar(segundos: int) -> bool:
        tempo_restante = segundos
        while tempo_restante > 0:
            sys.stdout.write('\r')
            tempo_restante -= 1
            sys.stdout.write(f"{tempo_restante}")
            sys.stdout.flush()
            time.sleep(1)
        return True


class OperacoesTarefa:
    """Operações de início e fim de tarefas."""
    
    @staticmethod
    def executar():
        """Executa a tarefa definida."""
        print("Executando tarefa")
        time.sleep(1)
    
    @staticmethod
    def finalizar():
        """Finaliza a tarefa."""
        print("Finalizando tarefa")


def worker_executor():
    """Executor de trabalho para o processo."""
    print('Iniciando worker')
    time.sleep(0.1)
    print('Worker concluído')
    OperacoesTarefa.executar()


def agendar_executar(horario: str = '', repeticoes: int = 1, timeout: int = 20):
    """Função de execução do agendamento."""
    if horario:
        ts_alvo = ConversorDataHora.converter_timestamp(horario)
        ts_atual = time.time()
        diff = int(ts_alvo - ts_atual)
        
        if diff < 0:
            raise ErroAgendador()
        
        while diff > 0:
            diff -= 1
            time.sleep(1)
    
    for _ in range(repeticoes):
        proc = multiprocessing.Process(target=worker_executor)
        proc.start()
        
        if timeout:
            result = ContadorRegressivo.iniciar(timeout)
            if result:
                proc.terminate()
                OperacoesTarefa.finalizar()


SINAL_PARADA = object()


class PoolThread:
    """Implementação de pool de threads para execução paralela."""
    
    def __init__(self, max_threads: int):
        self.fila = queue.Queue()
        self.max_threads = max_threads
        self.terminal = False
        self.threads_ativas = []
        self.threads_livres = []

    def submeter(
        self, 
        funcao: Callable, 
        args: Tuple, 
        callback: Optional[Callable] = None
    ) -> None:
        """Submete uma tarefa ao pool."""
        if not self.threads_livres and len(self.threads_ativas) < self.max_threads:
            self.criar_thread()
        
        trabalho = (funcao, args, callback)
        self.fila.put(trabalho)

    def criar_thread(self) -> None:
        """Cria uma nova thread worker."""
        t = threading.Thread(target=self.executar_loop)
        t.start()

    @contextlib.contextmanager
    def gerenciar_estado(self, lista: list, valor: Any):
        """Gerencia o estado da thread."""
        lista.append(valor)
        try:
            yield
        finally:
            lista.remove(valor)

    def executar_loop(self) -> None:
        """Loop principal de execução de tarefas."""
        thread_atual = threading.currentThread()
        self.threads_ativas.append(thread_atual)
        
        evento = self.fila.get()
        while evento != SINAL_PARADA:
            func, args, callback = evento
            
            try:
                resultado = func(*args)
                status = True
            except Exception as e:
                status = False
                resultado = e

            if callback:
                try:
                    callback(status, resultado)
                except Exception:
                    pass

            if self.terminal:
                evento = SINAL_PARADA
            else:
                with self.gerenciar_estado(self.threads_livres, thread_atual):
                    evento = self.fila.get()
        
        self.threads_ativas.remove(thread_atual)

    def encerrar(self) -> None:
        """Encerra o pool de threads."""
        qtd = len(self.threads_ativas)
        while qtd:
            self.fila.put(SINAL_PARADA)
            qtd -= 1

    def terminate(self) -> None:
        """Termina forçadamente o pool."""
        self.terminal = True
        while self.threads_ativas:
            self.fila.put(SINAL_PARADA)
        self.fila.empty()


if __name__ == '__main__':
    pool = PoolThread(20)

    for item in range(200):
        pool.submeter(func=agendar_executar, args=('', 1, 20))

    pool.terminate()

Exemplo de Uso Avançado

Podemos estender o sistema para cenários mais complexos:

#!/usr/bin/env python

import queue
import threading
import time
import os


class AgendadorTarefas:
    """Sistema de agendamento de tarefas avançado."""
    
    def __init__(self, max_concorrentes: int = 10):
        self.pool = PoolThread(max_concorrentes)
        self.em_execucao = True
    
    def executar_ciclo(self, tarefas: int, timeout: int):
        """Executa um ciclo de tarefas."""
        for _ in range(tarefas):
            self.pool.submeter(
                func=agendar_executar, 
                args=('', 1, timeout)
            )
    
    def iniciar(self, ciclos: int, tarefas_por_ciclo: int, timeout: int):
        """Inicia o agendador."""
        for ciclo in range(ciclos):
            self.executar_ciclo(tarefas_por_ciclo, timeout)
            time.sleep(2.5)
        
        self.pool.terminate()
        
        os.system("ps aux | grep python | grep -v grep | awk '{print $2}' | xargs -i kill -9 {}")


if __name__ == '__main__':
    agendador = AgendadorTarefas(10)
    agendador.iniciar(ciclos=100, tarefas_por_ciclo=3000, timeout=60)

Considerações Importantes

Ao implementar um sistema de agendamento, alguns pontos merecem atenção especial:

  • Gerenciamento de processos: Ao usar terminate(), o processo filho é encerrado, mas processos nipotes podem continuar em execução
  • Recursos do sistema: Cuidado com o número de processos/threads simultâneos para evitar sobrecarga
  • Tratamento de exceções: Sempre implemente tratamento adequado para evitar travamentos

Conclusão

A implementação de um sistema de agendamento em Python oferece flexibilidade equivalente ao Crontab do Linux, mas com maior controle sobre o ambiente de execução. As técnicas demonstradas permitem desde tarefas simples com timeout até processamento paralelo em larga escala.

Tags: Python Multiprocessing threading task-scheduler cron

Publicado em 10-1 11:25