0Pricing
Pandas & NumPy Academy · Aula

Agendando e registrando execuções do pipeline

Execute seu pipeline como um script Python pela linha de comando, registre os horários de início e término e use cron ou um agendador para automatizar a execução.

Agendando e registrando execuções do pipeline é uma aula grátis de Pandas & NumPy Academy no CoddyKit. Esta é a aula 4 de 4. Você pode ler a aula completa abaixo gratuitamente — depois pratica ao vivo no navegador com um editor de código integrado e um tutor de IA 24/7. Faz parte do caminho de aprendizado de Pandas & NumPy Academy, e seu progresso é sincronizado entre a web e o app CoddyKit. O curso de Pandas & NumPy Academy inclui 4 aulas no total.

Do notebook ao programa

Um fluxo de dados que é executado apenas quando um desenvolvedor abre manualmente um notebook não oferece valor comercial além da primeira execução. Para ser executado automaticamente todos os dias, o fluxo de dados precisa ser estruturado como um programa Python executável pela linha de comando: python pipeline.py. Isso exige um ponto de entrada if __name__ == '__main__':, análise dos argumentos da linha de comando e registro adequado — os três pilares de um programa de produção.

# pipeline.py
import argparse
import logging
import pandas as pd

def main(config_path):
    logging.info(f'Starting pipeline with config: {config_path}')
    # ... run ETL steps ...
    logging.info('Pipeline complete.')

if __name__ == '__main__':
    parser = argparse.ArgumentParser()
    parser.add_argument('--config', default='config.json')
    args = parser.parse_args()
    main(args.config)

Configurando o registro do Python

O módulo logging integrado ao Python é a ferramenta correta para os registros do fluxo de dados — não as instruções print(). Configure um registrador com saída no console e em arquivo usando logging.basicConfig(). Registre no nível INFO o progresso normal e no nível ERROR as falhas. Os registros baseados em arquivos permanecem disponíveis depois que o processo termina, o que é essencial para depurar execuções agendadas que ninguém estava acompanhando.

import logging
from datetime import date

log_file = f'pipeline_{date.today()}.log'

logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s %(levelname)s %(message)s',
    handlers=[
        logging.FileHandler(log_file),
        logging.StreamHandler()
    ]
)

logging.info('Logger configured.')

Registrando o início e o fim do fluxo de dados

Sempre registre a hora de início, a hora de término e o tempo decorrido de uma execução do fluxo de dados. Isso estabelece uma referência: se o fluxo de dados normalmente é executado em 45 segundos e hoje levou 8 minutos, algo mudou — talvez o arquivo de entrada esteja 10 vezes maior ou uma consulta ao banco de dados esteja lenta. As entradas de registro de início e término com carimbo de data e hora tornam essa comparação simples usando apenas o arquivo de registro.

import time
import logging

def run_pipeline(config):
    start = time.time()
    logging.info(f'Pipeline START | env={config.get("env", "dev")} | input={config["input_path"]}')

    try:
        df = extract(config)
        df_clean = transform(df, config)
        load(df_clean, config)
        elapsed = time.time() - start
        logging.info(f'Pipeline SUCCESS | rows={len(df_clean)} | elapsed={elapsed:.1f}s')
    except Exception as e:
        logging.error(f'Pipeline FAILED | error={e}', exc_info=True)
        raise

Registrando a quantidade de linhas por etapa

Registre a quantidade de linhas que entra e sai de cada etapa de transformação. Um registro claro se parece com: extract: 50,000 rows → drop_nulls: 49,200 rows → filter: 47,800 rows → output: 47,800 rows. Esse rastreamento deixa imediatamente claro quantas linhas foram removidas em cada etapa e se os números são esperados. Remoções anômalas aparecem como lacunas nas quantidades registradas.

def log_step(df, step_name):
    logging.info(f'{step_name}: {len(df):,} rows')
    return df

import pandas as pd

df = (pd.read_csv('orders.csv')
    .pipe(log_step, 'extract')
    .dropna(subset=['revenue'])
    .pipe(log_step, 'drop_nulls')
    .query('quantity > 0')
    .pipe(log_step, 'filter_qty')
)
print('Step logging complete.')

Agendando com cron no Linux/Mac

cron é o agendador Unix padrão para tarefas recorrentes. Edite o crontab com crontab -e e adicione uma linha especificando quando executar o programa. O formato é: minuto hora dia mês dia da semana comando. Um fluxo de dados que deve ser executado diariamente às 6h usa 0 6 * * * /usr/bin/python /path/to/pipeline.py. Sempre use caminhos absolutos nas entradas do cron, pois o cron é executado em um ambiente mínimo, sem as configurações de PATH do seu shell.

# crontab entry — edit with: crontab -e
# Run pipeline.py at 06:00 every day
# 0 6 * * * /opt/homebrew/bin/python /Users/analyst/pipeline.py --config /Users/analyst/config.json >> /Users/analyst/cron.log 2>&1

# Common cron patterns:
# 0 6 * * *      — daily at 06:00
# 0 */4 * * *    — every 4 hours
# 0 9 * * 1      — every Monday at 09:00
print('Cron schedule format: minute hour day month weekday')

Agendando com a biblioteca schedule do Python

A biblioteca schedule oferece uma forma totalmente baseada em Python de executar tarefas em intervalos especificados, sem usar o cron. Ela é útil em ambientes onde o cron não está disponível (Windows) ou quando você quer manter a lógica do agendador dentro do próprio processo Python. Envolva o fluxo de dados em um ciclo de tarefa agendada e mantenha o processo ativo para executá-lo repetidamente.

# pip install schedule
# import schedule, time

# def job():
#     logging.info('Scheduled run starting...')
#     run_pipeline(CONFIG)

# schedule.every().day.at('06:00').do(job)
# schedule.every(4).hours.do(job)

# while True:
#     schedule.run_pending()
#     time.sleep(60)

print('schedule library: use for in-process Python scheduling')

Tratamento de erros e códigos de saída

Um programa de fluxo de dados deve retornar um código de saída diferente de zero quando falhar, para que o agendador saiba que a tarefa falhou. Envolva a execução principal em um bloco try/except e chame sys.exit(1) em caso de falha. cron, Jenkins e Airflow verificam o código de saída: um código diferente de zero aciona um alerta, uma nova execução ou uma notificação. Uma exceção não tratada que não define o código de saída pode passar despercebida pelo monitoramento automatizado.

import sys

def main():
    try:
        run_pipeline(CONFIG)
        sys.exit(0)  # success
    except AssertionError as e:
        logging.error(f'Data validation failed: {e}')
        sys.exit(2)  # data error
    except Exception as e:
        logging.error(f'Unexpected error: {e}', exc_info=True)
        sys.exit(1)  # general failure

print('Exit code 0=success, 1=error, 2=data failure')

Escrevendo um arquivo de resumo da execução do fluxo de dados

Após uma execução bem-sucedida, escreva um pequeno arquivo de resumo JSON junto à saída. Inclua o carimbo de data e hora da execução, a quantidade de linhas de entrada, a quantidade de linhas de saída, as linhas removidas e o tempo decorrido. Sistemas de monitoramento e painéis podem ler esse arquivo para acompanhar as tendências de integridade do fluxo de dados ao longo do tempo. Um painel que mostre as linhas de saída nos últimos 30 dias facilita identificar o dia em que uma fonte de dados começou a fornecer menos registros.

import json
from datetime import datetime

def write_run_summary(config, input_rows, output_rows, elapsed):
    summary = {
        'run_at': datetime.now().isoformat(),
        'input_path': config['input_path'],
        'input_rows': input_rows,
        'output_rows': output_rows,
        'rows_dropped': input_rows - output_rows,
        'elapsed_seconds': round(elapsed, 2),
        'status': 'success'
    }
    with open('last_run_summary.json', 'w') as f:
        json.dump(summary, f, indent=2)
    print('Run summary written.')

Agendamento idempotente: evitando execuções duplicadas

Se um fluxo de dados agendado for acionado duas vezes por acidente, ele não deverá corromper a saída. Projete a etapa de carregamento para ser idempotente: use um nome de arquivo de saída com a data ou substitua a mesma saída pelo resultado mais recente. Para carregamentos em bancos de dados, use if_exists='replace' ou um padrão de UPSERT. Nunca use o modo append sem uma etapa de eliminação de duplicatas, pois cada execução agendada adicionará linhas duplicadas à tabela de saída.

from datetime import date

def load_idempotent(df, config):
    # Date-stamped output: each run overwrites its own day's file
    output_path = f"output_{date.today().strftime('%Y%m%d')}.parquet"
    df.to_parquet(output_path, index=False)
    logging.info(f'Loaded {len(df)} rows to {output_path}')

Alertas em caso de falha do fluxo de processamento

Para fluxos de processamento dos quais as operações da empresa dependem, o silêncio após uma falha é perigoso. Configure um alerta simples: se o arquivo de resumo da execução não for atualizado dentro do período esperado, envie um e-mail ou uma mensagem no Slack. O smtplib do Python pode enviar um e-mail em caso de falha, ou você pode usar um webhook para publicar uma mensagem no Slack. Gere um alerta imediatamente ao sair com o código 1 ou 2, para que o analista saiba que a atualização diária não foi concluída antes que a empresa perceba.

import smtplib

def send_failure_alert(error_msg):
    # Example: send plain-text email via SMTP
    # server = smtplib.SMTP('smtp.example.com', 587)
    # server.sendmail('pipeline@company.com',
    #                 'analyst@company.com',
    #                 f'Subject: Pipeline Failed\n\n{error_msg}')
    # server.quit()
    print(f'[ALERT] Would send failure notification: {error_msg}')

# In main():
# except Exception as e:
#     send_failure_alert(str(e))
#     sys.exit(1)
print('Alert integration pattern shown above.')

Script completo de fluxo de processamento agendado

Combine todas as partes — análise de argumentos, configuração de registro, resumo da execução, tratamento de erros e códigos de saída — em um script completo de fluxo de processamento. Esse script pode ser colocado em qualquer ambiente, apontado para um arquivo de configuração e agendado com o cron ou qualquer orquestrador de fluxos de trabalho. A cada execução, ele produz um arquivo de registro com data, um resumo da execução e um arquivo de saída com data, tornando cada execução totalmente auditável e reproduzível de forma independente.

# Full script skeleton:
# 1. parse --config argument
# 2. configure logging to file + console
# 3. load JSON config
# 4. validate config
# 5. run extract() -> transform() -> load()
# 6. write run summary JSON
# 7. sys.exit(0) on success, sys.exit(1) on failure

print('Production pipeline script structure complete.')
print('Schedule with: crontab -e  or  python scheduler.py')

Verificação rápida

Teste sua compreensão dos conceitos de Análise de Dados desta lição.

Recapitulação da lição

Nesta lição, você aprendeu: estruturar um fluxo de processamento como um script de linha de comando com análise de argumentos e registro, agendar com o cron e tratar falhas usando códigos de saída diferentes de zero e alertas e escrever arquivos de resumo da execução e projetar etapas de carregamento idempotentes para uma execução automatizada confiável. Parabéns por concluir a trilha de Análise de Dados: Pandas e NumPy!

Perguntas Frequentes

A aula “Agendando e registrando execuções do pipeline” é grátis?

Sim — o texto completo de “Agendando e registrando execuções do pipeline” é grátis para ler aqui na web. Para praticá-la interativamente (um editor de código integrado e um tutor de IA 24/7) e desbloquear o restante do curso de Pandas & NumPy Academy, atualize para CoddyKit PRO. O curso de Pandas & NumPy Academy inclui 4 aulas no total.

O que vou aprender em “Agendando e registrando execuções do pipeline”?

Execute seu pipeline como um script Python pela linha de comando, registre os horários de início e término e use cron ou um agendador para automatizar a execução. Você pratica Pandas & NumPy Academy com código prático que executa diretamente no navegador, e um tutor de IA 24/7 responde suas dúvidas enquanto trabalha na aula.

Preciso ter experiência prévia para começar Pandas & NumPy Academy?

Nenhuma experiência prévia é necessária. Pandas & NumPy Academy no CoddyKit é estruturado para alunos iniciantes até avançados, então você pode começar aqui ou desde o início e aprender no seu ritmo. Esta é a aula 4 de 4.

Quanto tempo leva a aula “Agendando e registrando execuções do pipeline”?

A maioria das aulas CoddyKit leva cerca de 5–10 minutos. Cada uma é compacta e interativa, então você faz progresso constante e retoma exatamente de onde parou entre web e app.

Posso escrever e executar código nesta aula de Pandas & NumPy Academy?

Sim. Cada aula de Pandas & NumPy Academy inclui um editor de código integrado, então você escreve e executa código real direto no navegador e recebe feedback de IA instantaneamente — nenhuma configuração local necessária.

Todas as aulas deste curso

  1. Estruturando etapas de transformação como funções
  2. Parametrizando pipelines com dicionários de configuração
  3. Testando etapas do pipeline com asserções
  4. Agendando e registrando execuções do pipeline
← Voltar para Pandas & NumPy Academy