Trigger spark job using Airflow

Viewed 341
  • I am working on airflow and Spark in two separate containers.

  • I set up a spark-master and spark-worker1 connected through a bridge network called Spark_net.

  • Airflow runs using docker compose and has its own bridge network called airflow_default.

  • I connected the tow spark containers to the network airflow default using docker network connect airflow_default spark-worker1 and docker network connect airflow_default spark-master .

  • I think till this point the setup is correct.

  • Now, I would like to trigger a spark job using airflow :

    • in airflow, I added a new connection " Admin > Connections > Add new record "enter image description here
      • connection Id = spark_default
      • connection type = Spark
      • host = here I did docker network inspect airflow_default to get the Ip-address of the spark-master
    • I used SparkSubmitOperator to submit the job ( the job is just print a dataframe )

the DAG

from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from pyspark.sql import SparkSession
from datetime import datetime, date
import pandas as pd
from pyspark.sql import Row
    
default_args = {'owner': 'airflow','start_date': datetime(2021, 5, 9),}
    
dag = DAG('spark', description = 'spark_test', catchup = False, schedule_interval = None, default_args = default_args)

spark_job = SparkSubmitOperator(task_id = "spark_job",
                                application = "./data/spark-app.py",
                                conn_id = "spark_default",
                                dag = dag)
    
spark_job

the spark-app.py

from pyspark import SparkContext

sc = SparkContext("local", "First App")
    
data = [{"Category": 'A', "ID": 1, "Value": 121.44, "Truth": True},
        {"Category": 'B', "ID": 2, "Value": 300.01, "Truth": False},
        {"Category": 'C', "ID": 3, "Value": 10.99, "Truth": None},
        {"Category": 'E', "ID": 4, "Value": 33.87, "Truth": True}]
    
df = sc.parallelize(data)
df = df.collect()

print(df)

Now the problem is when I run the DAG, it works perfectly but in the spark-master UI, there is no job running or worker executing something.

I thought airflow is not communicating with spark-master, So I tried changing the IP-address in the connection and I put nonsense IP@ like 127.127.124.127 and the DAG still works perfectly and displays the dataframe

How can fix the issue and really connect to spark-master and trigger job from airflow that really are passed to the spark-master I have

0 Answers
Related