Harmonogramowanie i rejestrowanie uruchomień potoku
Uruchamiaj potok jako skrypt Pythona z wiersza poleceń, rejestruj czas rozpoczęcia i zakończenia oraz używaj cron lub harmonogramu do automatyzacji.
Harmonogramowanie i rejestrowanie uruchomień potoku to bezpłatna lekcja Pandas & NumPy Academy na CoddyKit. To lekcja 4 z 4. Możesz przeczytać całą lekcję poniżej za darmo — a potem ćwiczyć ją interaktywnie w przeglądarce z wbudowanym edytorem kodu i tutorem AI dostępnym 24/7. To część ścieżki edukacyjnej Pandas & NumPy Academy, a Twój postęp synchronizuje się między webem a aplikacją CoddyKit. Kurs Pandas & NumPy Academy zawiera 4 lekcji w sumie.
Od notebooka do skryptu
Pipeline uruchamiany tylko wtedy, gdy deweloper ręcznie otworzy notebook, nie przynosi żadnej wartości biznesowej poza pierwszym uruchomieniem. Aby uruchamiać go automatycznie każdego dnia, należy zorganizować pipeline jako wykonywalny z wiersza poleceń skrypt Pythona script: python pipeline.py. Wymaga to punktu wejścia if __name__ == '__main__':, analizowania argumentów wiersza poleceń oraz poprawnego rejestrowania zdarzeń — są to trzy filary skryptu produkcyjnego.
# 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)Konfigurowanie rejestrowania zdarzeń w Pythonie
Wbudowany w Pythona moduł logging jest właściwym narzędziem do rejestrowania logów pipeline’u — nie instrukcje print(). Należy skonfigurować logger z wyjściem zarówno na konsolę, jak i do pliku, używając logging.basicConfig(). Do rejestrowania zwykłego postępu należy używać poziomu INFO, a do błędów — ERROR. Logi zapisywane w pliku pozostają dostępne po zakończeniu procesu, co jest niezbędne do debugowania zaplanowanych uruchomień, których nikt nie obserwował.
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.')Rejestrowanie początku i końca pipeline’u
Zawsze należy rejestrować czas rozpoczęcia, czas zakończenia i czas trwania uruchomienia pipeline’u. Tworzy to punkt odniesienia: jeśli pipeline zwykle działa przez 45 sekund, a dziś działał przez 8 minut, coś się zmieniło — być może plik wejściowy jest 10 razy większy albo zapytanie do bazy danych działa wolniej. Wpisy dziennika z oznaczeniem czasu, informujące o rozpoczęciu i zakończeniu, pozwalają łatwo porównać te wartości na podstawie samego pliku logu.
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)
raiseRejestrowanie liczby wierszy na każdym kroku
Należy rejestrować liczbę wierszy wchodzących do każdego kroku transformacji i wychodzących z niego. Przejrzysty log może wyglądać tak: extract: 50,000 rows → drop_nulls: 49,200 rows → filter: 47,800 rows → output: 47,800 rows. Taki ślad od razu pokazuje, ile wierszy usunięto na każdym kroku i czy liczby są zgodne z oczekiwaniami. Nietypowe spadki są widoczne jako luki między zarejestrowanymi wartościami.
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.')Planowanie za pomocą cron w systemach Linux/Mac
cron to standardowy harmonogram zadań cyklicznych w systemach Unix. Należy edytować crontab za pomocą crontab -e i dodać wiersz określający, kiedy uruchomić skrypt. Format to: minuta godzina dzień miesiąc dzień_tygodnia polecenie. Pipeline, który musi uruchamiać się codziennie o 6:00 rano, używa wpisu 0 6 * * * /usr/bin/python /path/to/pipeline.py. Wpisy cron powinny zawsze zawierać ścieżki absolutne, ponieważ cron działa w minimalnym środowisku, bez ustawień PATH powłoki.
# 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')Planowanie za pomocą biblioteki Python schedule
Biblioteka schedule zapewnia sposób uruchamiania zadań w określonych odstępach czasu bez używania cron, napisany w całości w Pythonie. Jest przydatna w środowiskach, w których cron nie jest dostępny (na przykład w systemie Windows), lub gdy logika harmonogramu ma znajdować się bezpośrednio w procesie Pythona. Należy opakować pipeline w pętlę zaplanowanych zadań i utrzymywać proces uruchomiony, aby zadania wykonywały się wielokrotnie.
# 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')Obsługa błędów i kody wyjścia
Skrypt pipeline’u powinien zwracać niezerowy kod wyjścia w przypadku awarii, aby harmonogram wiedział, że zadanie zakończyło się niepowodzeniem. Główne wykonanie należy opakować w blok try/except i w razie błędu wywołać sys.exit(1). cron, Jenkins i Airflow sprawdzają kod wyjścia: kod niezerowy uruchamia alert, ponowienie zadania albo powiadomienie. Nieobsłużony wyjątek, który nie ustawia kodu wyjścia, może pozostać niezauważony przez automatyczne monitorowanie.
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')Tworzenie pliku podsumowania uruchomienia pipeline’u
Po pomyślnym uruchomieniu należy zapisać obok wyniku niewielki plik podsumowania w formacie JSON. Należy uwzględnić w nim znacznik czasu uruchomienia, liczbę wierszy wejściowych, liczbę wierszy wynikowych, liczbę usuniętych wierszy oraz czas trwania. Systemy monitorowania i pulpity nawigacyjne mogą odczytywać ten plik, aby śledzić trendy dotyczące kondycji pipeline’u w czasie. Pulpit pokazujący liczbę wierszy wynikowych z ostatnich 30 dni ułatwia wskazanie dnia, w którym źródło danych zaczęło dostarczać mniej rekordów.
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.')Idempotentne planowanie: unikanie podwójnych uruchomień
Jeśli zaplanowany pipeline zostanie przypadkowo uruchomiony dwa razy, nie powinien uszkodzić wyniku. Krok ładowania należy zaprojektować jako idempotentny: można używać nazwy pliku wynikowego zawierającej datę albo nadpisywać ten sam plik najnowszym wynikiem. W przypadku ładowania do bazy danych należy użyć if_exists='replace' albo wzorca UPSERT. Nie należy nigdy używać trybu append bez kroku usuwania duplikatów, ponieważ każde zaplanowane uruchomienie doda zduplikowane wiersze do tabeli wynikowej.
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}')Alertowanie o awarii potoku
W przypadku potoków, od których zależy działanie firmy, brak informacji po awarii jest niebezpieczny. Proszę skonfigurować prosty alert: jeśli plik podsumowania uruchomienia nie zostanie zaktualizowany w oczekiwanym czasie, należy wysłać wiadomość e-mail lub w Slacku. smtplib w Pythonie może wysłać wiadomość e-mail w razie awarii. Można też użyć webhooka do opublikowania wiadomości w Slacku. Należy natychmiast generować alert przy kodzie wyjścia 1 lub 2, aby analityk wiedział o nieudanym codziennym odświeżeniu, zanim zauważy to firma.
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.')Kompletny skrypt planowanego potoku
Proszę połączyć wszystkie elementy — analizowanie argumentów, konfigurację logowania, podsumowanie uruchomienia, obsługę błędów i kody wyjścia — w kompletny skrypt potoku. Taki skrypt można wdrożyć w dowolnym środowisku, wskazać mu plik konfiguracyjny i zaplanować jego uruchamianie za pomocą cron lub dowolnego orkiestratora przepływów pracy. Przy każdym uruchomieniu tworzy on plik dziennika z datą, podsumowanie uruchomienia oraz plik wyjściowy z datą, dzięki czemu każde uruchomienie można w pełni skontrolować i niezależnie odtworzyć.
# 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')Szybki test
Sprawdź swoją znajomość zagadnień analizy danych omówionych w tej lekcji.
Podsumowanie lekcji
W tej lekcji nauczyłeś się: tworzenia struktury potoku jako skryptu wiersza poleceń z analizowaniem argumentów i logowaniem, planowania zadań za pomocą cron oraz obsługi awarii z użyciem niezerowych kodów wyjścia i alertów, a także zapisywania plików podsumowania uruchomienia i projektowania idempotentnych etapów ładowania na potrzeby niezawodnego zautomatyzowanego wykonywania. Gratulacje z okazji ukończenia ścieżki Analiza danych: Pandas i NumPy!
Ucz się Python dzięki korepetycjom AI — za darmo
Pisz i uruchamiaj kod w przeglądarce, otrzymuj natychmiastową pomoc od korepetytora AI dostępnego 24/7 i kontynuuj naukę w sieci lub w aplikacji.
- Kursy
- 30
- Lekcje
- 120
Często zadawane pytania
Czy lekcja „Harmonogramowanie i rejestrowanie uruchomień potoku” jest bezpłatna?
Tak — pełny tekst „Harmonogramowanie i rejestrowanie uruchomień potoku” jest dostępny za darmo tutaj w sieci. Aby ćwiczyć ją interaktywnie (wbudowany edytor kodu i tutor AI dostępny 24/7) i odblokować resztę kursu Pandas & NumPy Academy, przejdź na CoddyKit PRO. Kurs Pandas & NumPy Academy zawiera 4 lekcji w sumie.
Co nauczysz się w „Harmonogramowanie i rejestrowanie uruchomień potoku”?
Uruchamiaj potok jako skrypt Pythona z wiersza poleceń, rejestruj czas rozpoczęcia i zakończenia oraz używaj cron lub harmonogramu do automatyzacji. Ćwiczysz Pandas & NumPy Academy z praktycznym kodem, który uruchamiasz bezpośrednio w przeglądarce, a tutor AI dostępny 24/7 odpowiada na Twoje pytania podczas pracy nad lekcją.
Czy potrzebuję doświadczenia, aby zacząć Pandas & NumPy Academy?
Nie wymagamy żadnego doświadczenia. Pandas & NumPy Academy w CoddyKit jest strukturyzowany dla początkujących i zaawansowanych użytkowników, więc możesz zacząć tutaj lub od początku i uczyć się w swoim tempie. To lekcja 4 z 4.
Ile czasu zajmuje lekcja „Harmonogramowanie i rejestrowanie uruchomień potoku”?
Większość lekcji CoddyKit trwa około 5–10 minut. Każda lekcja to mały, interaktywny krok, dzięki czemu robisz systematyczne postępy i zawsze wracasz dokładnie do tego samego miejsca — na webie i w aplikacji.
Czy mogę pisać i uruchamiać kod w tej lekcji Pandas & NumPy Academy?
Tak. Każda lekcja Pandas & NumPy Academy zawiera wbudowany edytor kodu, więc piszesz i uruchamiasz prawdziwy kod bezpośrednio w przeglądarce i od razu otrzymujesz sprzężenie zwrotne od AI — bez konfiguracji na komputerze.
Wszystkie lekcje w tym kursie
- Strukturyzowanie etapów transformacji jako funkcji
- Parametryzowanie potoków za pomocą słowników konfiguracji
- Testowanie etapów potoku za pomocą asercji
- Harmonogramowanie i rejestrowanie uruchomień potoku