Pandas & NumPy Academy · Урок

Планирование и журналирование запусков конвейера

Запускайте конвейер как скрипт Python из командной строки, записывайте время начала и окончания и используйте cron или планировщик для автоматизации.

Урок 4 из 413 шагов

«Планирование и журналирование запусков конвейера» — бесплатный урок Pandas & NumPy Academy на CoddyKit. Это урок 4 из 4. Ты можешь прочитать весь урок бесплатно ниже — а потом практиковать его прямо в браузере с встроенным редактором кода и ИИ-репетитором 24/7. Это часть пути обучения Pandas & NumPy Academy, и твой прогресс синхронизируется между веб-версией и приложением CoddyKit. Курс Pandas & NumPy Academy содержит 4 уроков всего.

От блокнота к скрипту

Конвейер, который запускается только тогда, когда разработчик вручную открывает блокнот, не приносит бизнес-ценности после первого запуска. Чтобы запускать его автоматически каждый день, необходимо оформить конвейер как исполняемый из командной строки Python-скрипт: python pipeline.py. Для этого нужны точка входа if __name__ == '__main__':, разбор аргументов командной строки и корректное журналирование — три основы рабочего скрипта.

# 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)

Настройка журналирования в Python

Встроенный модуль Python logging — правильный инструмент для журналов конвейера, а не операторы print(). Настройте регистратор с выводом и в консоль, и в файл с помощью logging.basicConfig(). Используйте уровень INFO для обычного хода выполнения и ERROR для сбоев. Журналы в файлах сохраняются после завершения процесса, что особенно важно для отладки запланированных запусков, за которыми никто не наблюдал.

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.')

Журналирование начала и окончания конвейера

Всегда записывайте время начала, время окончания и продолжительность запуска конвейера. Это создаёт базовую точку для сравнения: если обычно конвейер выполняется за 45 секунд, а сегодня работал 8 минут, значит, что-то изменилось — например, входной файл стал в 10 раз больше или запрос к базе данных выполняется медленно. Записи о начале и окончании с временными метками позволяют легко сравнить эти показатели, просматривая только файл журнала.

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

Журналирование количества строк на каждом шаге

Записывайте количество строк до и после каждого шага преобразования. Чистый журнал может выглядеть так: извлечение: 50 000 строк → drop_nulls: 49 200 строк → фильтрация: 47 800 строк → результат: 47 800 строк. По такой трассировке сразу видно, сколько строк было удалено на каждом шаге и соответствуют ли значения ожиданиям. Аномальные уменьшения проявляются как разрывы в записанных значениях.

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.')

Планирование с помощью cron в Linux и Mac

cron — стандартный планировщик Unix для регулярно выполняемых заданий. Отредактируйте таблицу заданий с помощью crontab -e и добавьте строку, указывающую время запуска скрипта. Формат: минута час день месяц день_недели команда. Для ежедневного запуска конвейера в 6:00 утра используется 0 6 * * * /usr/bin/python /path/to/pipeline.py. Всегда указывайте в записях cron абсолютные пути, поскольку cron работает в минимальном окружении без настроек PATH вашей оболочки.

# 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')

Планирование с помощью библиотеки schedule для Python

Библиотека schedule позволяет запускать задания через заданные интервалы средствами чистого Python, не обращаясь к cron. Она полезна в средах, где cron недоступен (например, в Windows), или когда логика планировщика должна находиться непосредственно в процессе Python. Оберните конвейер в цикл запланированного задания и поддерживайте процесс работающим, чтобы задания выполнялись повторно.

# 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')

Обработка ошибок и коды завершения

При сбое скрипт конвейера должен возвращать ненулевой код завершения, чтобы планировщик знал о неудаче задания. Оберните основное выполнение в блок try/except и при сбое вызовите sys.exit(1). cron, Jenkins и Airflow проверяют код завершения: ненулевой код запускает оповещение, повторный запуск или уведомление. Необработанное исключение, не устанавливающее код завершения, может остаться незамеченным автоматическим мониторингом.

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')

Запись файла со сводкой запуска конвейера

После успешного запуска записывайте небольшой файл JSON со сводкой рядом с результатом. Включите в него временную метку запуска, количество входных и выходных строк, число удалённых строк и продолжительность выполнения. Системы мониторинга и панели мониторинга могут читать этот файл, чтобы отслеживать изменения состояния конвейера во времени. Панель с показателем количество выходных строк за последние 30 дней позволяет легко заметить день, когда источник данных начал передавать меньше записей.

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.')

Идемпотентное планирование: предотвращение повторных запусков

Если запланированный конвейер случайно запустился дважды, это не должно повредить результат. Спроектируйте шаг загрузки как идемпотентный: используйте имя выходного файла с датой или перезаписывайте один и тот же результат последней версией. Для загрузки в базу данных используйте if_exists='replace' или шаблон UPSERT. Никогда не используйте режим append без шага удаления дубликатов, иначе каждый запланированный запуск будет добавлять повторяющиеся строки в выходную таблицу.

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}')

Оповещение о сбое конвейера

Для конвейеров, от которых зависят рабочие процессы компании, молчание после сбоя опасно. Настройте простое оповещение: если файл со сводкой запуска не обновился в ожидаемый срок, отправляйте электронное письмо или сообщение в Slack. smtplib в Python может отправить письмо при сбое, либо можно использовать веб-перехватчик для публикации сообщения в Slack. Немедленно оповещайте о коде завершения 1 или 2, чтобы аналитик узнал о пропущенном ежедневном обновлении раньше бизнеса.

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.')

Полный скрипт запланированного конвейера

Объедините все части — разбор аргументов, настройку журналирования, сводку запуска, обработку ошибок и коды завершения — в полный скрипт конвейера. Такой скрипт можно перенести в любую среду, указать ему файл конфигурации и запланировать с помощью cron или любого оркестратора рабочих процессов. При каждом запуске он создаёт файл журнала с датой, сводку запуска и выходной файл с датой, благодаря чему каждый запуск полностью проверяем и воспроизводим независимо.

# 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')

Быстрая проверка

Проверьте, насколько хорошо Вы усвоили концепции анализа данных из этого урока.

Итоги урока

В этом уроке Вы узнали, как структурировать конвейер в виде скрипта командной строки с разбором аргументов и журналированием, планировать запуски с помощью cron и обрабатывать сбои с ненулевыми кодами завершения и оповещениями, а также создавать файлы со сводками запусков и проектировать идемпотентные шаги загрузки для надёжного автоматического выполнения. Поздравляем с завершением направления «Анализ данных: Pandas и NumPy»!

Можно начать бесплатно

Изучай Python с ИИ-репетитором — бесплатно

Пиши и запускай код прямо в браузере, получай мгновенную помощь от ИИ-репетитора 24/7 и продолжи учиться на сайте или в приложении.

Курсы
30
Уроки
120

Часто задаваемые вопросы

Урок «Планирование и журналирование запусков конвейера» бесплатный?

Да — полный текст урока «Планирование и журналирование запусков конвейера» бесплатно доступен здесь в веб-версии. Чтобы практиковать его интерактивно (встроенный редактор кода и ИИ-репетитор 24/7) и разблокировать остальной курс Pandas & NumPy Academy, подпишись на CoddyKit PRO. Курс Pandas & NumPy Academy содержит 4 уроков всего.

Чему я научусь в уроке «Планирование и журналирование запусков конвейера»?

Запускайте конвейер как скрипт Python из командной строки, записывайте время начала и окончания и используйте cron или планировщик для автоматизации. Ты практикуешь Pandas & NumPy Academy с помощью реального кода, который запускаешь прямо в браузере, и ИИ-репетитор 24/7 отвечает на твои вопросы во время урока.

Нужен ли мне опыт, чтобы начать Pandas & NumPy Academy?

Предыдущий опыт не требуется. Pandas & NumPy Academy на CoddyKit структурирован для всех уровней — от новичков до продвинутых, поэтому ты можешь начать отсюда или с самого начала и учиться в своем темпе. Это урок 4 из 4.

Сколько времени занимает урок «Планирование и журналирование запусков конвейера»?

Большинство уроков CoddyKit занимают около 5–10 минут. Каждый из них компактный и интерактивный, поэтому ты постоянно делаешь прогресс и продолжаешь с того же места в веб-версии и приложении.

Можно ли писать и запускать код в этом уроке Pandas & NumPy Academy?

Да. Каждый урок Pandas & NumPy Academy включает встроенный редактор кода, поэтому ты пишешь и запускаешь реальный код прямо в браузере и получаешь моментальную обратную связь от AI — локальная установка не требуется.

Все уроки этого курса

  1. Оформление шагов преобразования в виде функций
  2. Параметризация конвейеров с помощью словарей конфигурации
  3. Тестирование этапов конвейера с помощью утверждений
  4. Планирование и журналирование запусков конвейера
← Назад к Pandas & NumPy Academy