パイプライン実行のスケジューリングとログ記録
コマンドラインからPythonスクリプトとしてパイプラインを実行し、開始時刻と終了時刻を記録して、cronまたはスケジューラーで自動化します。
「パイプライン実行のスケジューリングとログ記録」はCoddyKit上の無料Pandas & NumPy Academyレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはPandas & NumPy Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Pandas & NumPy Academyコースには全4レッスンが含まれています。
ノートブックからスクリプトへ
開発者が手動でノートブックを開いたときにしか実行されないパイプラインは、最初の実行以外にはビジネス上の価値を提供しません。毎日自動的に実行するには、パイプラインをコマンドラインから実行できる Python のスクリプトとして構成する必要があります:python pipeline.py。そのためには、if __name__ == '__main__': のエントリーポイント、コマンドライン引数の解析、適切なログ記録という、プロダクションスクリプトの 3 つの柱が必要です。
# 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 ロギングの設定
パイプラインのログには、print() 文ではなく、Python 組み込みの logging モジュールを使うのが適切です。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各ステップの行数をログに記録する
各変換ステップに入る行数と、そこから出る行数を記録します。整ったログは次のようになります:extract: 50,000 行 → drop_nulls: 49,200 行 → filter: 47,800 行 → output: 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.')Linux/Mac で cron を使ってスケジュールする
cron は、定期的なジョブを実行するための標準的な Unix スケジューラーです。crontab -e で crontab を編集し、スクリプトを実行する時刻を指定する行を追加します。形式は次のとおりです:minute hour day month weekday command。毎日午前 6:00 に実行するパイプラインには、0 6 * * * /usr/bin/python /path/to/pipeline.py を使用します。cron はシェルの PATH 設定がない最小限の環境で実行されるため、cron のエントリでは必ず絶対パスを使用します。
# 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')Python schedule ライブラリでスケジュールする
schedule ライブラリを使うと、cron に触れずに、指定した間隔でジョブを実行する純粋な Python の方法を利用できます。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')エラー処理と終了コード
パイプラインスクリプトは、失敗したときに 0 以外の終了コードを返し、スケジューラーがジョブの失敗を認識できるようにする必要があります。メイン処理を try/except ブロックで囲み、失敗時に sys.exit(1) を呼び出します。cron、Jenkins、Airflow はいずれも終了コードを確認します。0 以外のコードによって、アラート、再実行、通知が発生します。終了コードを設定しない未処理の例外は、自動監視で見逃される可能性があります。
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.')冪等なスケジューリング:二重実行を防ぐ
スケジュールされたパイプラインが誤って 2 回起動されても、出力が壊れないようにする必要があります。日付付きの出力ファイル名を使う、または最新の結果で同じ出力を上書きするなど、load ステップを冪等に設計します。データベースへのロードでは、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メッセージを送信します。Pythonの smtplib を使えば失敗時にメールを送信できます。また、Webhookを使って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')クイックチェック
このレッスンで学んだData Analysisの概念について理解度を確認します。
レッスンのまとめ
このレッスンでは、引数解析とログ機能を備えたコマンドラインスクリプトとしてパイプラインを構成する方法、cronでスケジュール実行し、ゼロ以外の終了コードとアラートで失敗を処理する方法、そして信頼性の高い自動実行のために実行サマリーファイルを作成し、冪等なロード処理を設計する方法を学びました。Data Analysis: Pandas and NumPyトラックの修了、おめでとうございます。
AI チューターと学ぶ Python — 無料
ブラウザでリアルコードを書いて実行し、24/7 の AI チューターから瞬時にサポートを受け、ウェブまたはアプリで続きから学習できます。
- コース
- 30
- レッスン
- 120
よくある質問
「パイプライン実行のスケジューリングとログ記録」レッスンは無料ですか?
はい。「パイプライン実行のスケジューリングとログ記録」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Pandas & NumPy Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Pandas & NumPy Academyコースには全4レッスンが含まれています。
「パイプライン実行のスケジューリングとログ記録」で何を学びますか?
コマンドラインからPythonスクリプトとしてパイプラインを実行し、開始時刻と終了時刻を記録して、cronまたはスケジューラーで自動化します。 ブラウザで直接実行するハンズオンコードでPandas & NumPy Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
Pandas & NumPy Academyを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのPandas & NumPy Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン4/4です。
「パイプライン実行のスケジューリングとログ記録」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このPandas & NumPy Academyレッスンでコードを書いて実行できますか?
はい。すべてのPandas & NumPy Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- 変換手順の関数化
- 設定用dictによるパイプラインのパラメーター化
- アサーションによるパイプライン手順のテスト
- パイプライン実行のスケジューリングとログ記録