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.