Frameworki do harmonogramowania przepływów pracy

Wprowadzenie do inżynierii danych

Vincent Vankrunkelsven

Data Engineer @ DataCamp

Przykładowy potok danych

 

Przykładowy prosty potok danych wyodrębniający dane z CSV za pomocą Spark

Jak harmonogramować?

  • Ręcznie
  • Narzędzie cron
  • A zależności?
Wprowadzenie do inżynierii danych

DAG-i

Directed Acyclic Graph

  • Zbiór węzłów
  • Krawędzie skierowane
  • Brak cykli

Przykładowy DAG

Wprowadzenie do inżynierii danych

Narzędzia do zadania

 

  • cron w systemie Linux
  • Luigi od Spotify
  • Apache Airflow
Wprowadzenie do inżynierii danych

Logo Apache Airflow

  • Stworzony w Airbnb
  • DAG-i
  • Python
Wprowadzenie do inżynierii danych

Airflow: przykładowy DAG

 

Przykładowy DAG w Airflow

Wprowadzenie do inżynierii danych

Airflow: przykład w kodzie

# Create the DAG object
dag = DAG(dag_id="example_dag", ..., schedule_interval="0 * * * *")

# Define operations start_cluster = StartClusterOperator(task_id="start_cluster", dag=dag) ingest_customer_data = SparkJobOperator(task_id="ingest_customer_data", dag=dag) ingest_product_data = SparkJobOperator(task_id="ingest_product_data", dag=dag) enrich_customer_data = PythonOperator(task_id="enrich_customer_data", ..., dag = dag)
# Set up dependency flow start_cluster.set_downstream(ingest_customer_data) ingest_customer_data.set_downstream(enrich_customer_data) ingest_product_data.set_downstream(enrich_customer_data)
Wprowadzenie do inżynierii danych

Czas na ćwiczenia!

Wprowadzenie do inżynierii danych

Preparing Video For Download...