How to work with Apache Airflow Data in Apache Spark using SQL
Apache Spark is a fast and general engine for large-scale data processing. When paired with the CData JDBC Driver for Apache Airflow, Spark can work with live Apache Airflow data. This article describes how to connect to and query Apache Airflow data from a Spark shell.
The CData JDBC Driver offers unmatched performance for interacting with live Apache Airflow data due to optimized data processing built into the driver. When you issue complex SQL queries to Apache Airflow, the driver pushes supported SQL operations, like filters and aggregations, directly to Apache Airflow and utilizes the embedded SQL engine to process unsupported operations (often SQL functions and JOIN operations) client-side. With built-in dynamic metadata querying, you can work with and analyze Apache Airflow data using native data types.
Install the CData JDBC Driver for Apache Airflow
Download the CData JDBC Driver for Apache Airflow installer, unzip the package, and run the JAR file to install the driver.
Start a Spark Shell and Connect to Apache Airflow Data
- Open a terminal and start the Spark shell with the CData JDBC Driver for Apache Airflow JAR file as the jars parameter:
$ spark-shell --jars /CData/CData JDBC Driver for Apache Airflow/lib/cdata.jdbc.api.jar - With the shell running, you can connect to Apache Airflow with a JDBC URL and use the SQL Context load() function to read a table.
Start by setting the Profile connection property to the location of the ApacheAirflow Profile on disk (e.g. C:\profiles\ApacheAirflow.apip). Next, set the ProfileSettings connection property to the connection string for ApacheAirflow (see below).
ApacheAirflow API Profile Settings
Apache Airflow 3 uses JWT Bearer tokens for API authentication. You can generate a token from the Airflow web UI under Settings or via the Airflow CLI using airflow users create and the /api/v2/auth/token endpoint. Note that this profile targets Apache Airflow 3.x using the /api/v2 REST API. The legacy /api/v1 endpoint used by Airflow 2.x is not supported.
After setting the following connection properties, you are ready to connect:
- AuthScheme: Set this to APIKey.
- APIKey: Set this to your Apache Airflow JWT Bearer token.
- Server: Set this to the base URL of your Airflow instance (e.g. http://localhost:8080).
Built-in Connection String Designer
For assistance in constructing the JDBC URL, use the connection string designer built into the Apache Airflow JDBC Driver. Either double-click the JAR file or execute the jar file from the command-line.
java -jar cdata.jdbc.api.jarFill in the connection properties and copy the connection string to the clipboard.
Configure the connection to Apache Airflow, using the connection string generated above.
scala> val api_df = spark.sqlContext.read.format("jdbc").option("url", "jdbc:api:Profile=C:\profiles\ApacheAirflow.apip;AuthScheme=APIKey;ProfileSettings='APIKey=your_jwt_token;Server=http://localhost:8080';").option("dbtable","DagRuns").option("driver","cdata.jdbc.api.APIDriver").load() - Once you connect and the data is loaded you will see the table schema displayed.
Register the Apache Airflow data as a temporary table:
scala> api_df.registerTable("dagruns")-
Perform custom SQL queries against the Data using commands like the one below:
scala> api_df.sqlContext.sql("SELECT DagRunId, State FROM DagRuns WHERE DagId = example_dag").collect.foreach(println)You will see the results displayed in the console, similar to the following:
Using the CData JDBC Driver for Apache Airflow in Apache Spark, you are able to perform fast and complex analytics on Apache Airflow data, combining the power and utility of Spark with your data. Download a free, 30 day trial of any of the hundreds of CData JDBC Drivers and get started today.