Планирование и журналирование запусков конвейера
Запускайте конвейер как скрипт Python из командной строки, записывайте время начала и окончания и используйте cron или планировщик для автоматизации.
«Планирование и журналирование запусков конвейера» — бесплатный урок 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 — локальная установка не требуется.
Все уроки этого курса
- Оформление шагов преобразования в виде функций
- Параметризация конвейеров с помощью словарей конфигурации
- Тестирование этапов конвейера с помощью утверждений
- Планирование и журналирование запусков конвейера