如何通过Airflow Connection使用PySpark查询MySQL并加载为DataFrame?
使用Airflow Connection将MySQL数据加载到PySpark DataFrame
当然可以通过Airflow Connection实现——Airflow Connection存储的MySQL连接凭据可直接提取,作为PySpark JDBC连接的参数,既避免硬编码敏感信息,又能统一管理连接配置。
实现步骤与代码示例
1. 提取Airflow MySQL连接信息
用Airflow的BaseHook获取已配置的MySQL连接对象,从中提取所需参数:
from airflow.hooks.base import BaseHook from pyspark.sql import SparkSession # 替换为你在Airflow中配置的MySQL连接ID MYSQL_CONN_ID = "mysql_prod_conn" # 获取Airflow连接对象 conn = BaseHook.get_connection(MYSQL_CONN_ID) # 提取核心连接参数 mysql_host = conn.host mysql_port = conn.port mysql_user = conn.login mysql_password = conn.password mysql_db = conn.schema
2. 构建PySpark JDBC配置
将提取的参数组装成JDBC URL和连接属性:
# 构建MySQL JDBC URL(可根据需求调整SSL、时区参数) jdbc_url = f"jdbc:mysql://{mysql_host}:{mysql_port}/{mysql_db}?useSSL=false&serverTimezone=UTC" # JDBC连接属性 jdbc_properties = { "user": mysql_user, "password": mysql_password, "driver": "com.mysql.cj.jdbc.Driver" # MySQL 8.0+推荐使用的驱动 }
3. 加载MySQL数据到PySpark DataFrame
初始化SparkSession后,通过read.jdbc()方法加载数据:
# 初始化SparkSession(需指定JDBC驱动路径) spark = SparkSession.builder \ .appName("LoadMySQLDataViaAirflow") \ .config("spark.jars", "mysql-connector-java-8.0.30.jar") # 替换为你的JDBC驱动路径/版本 .getOrCreate() # 加载数据(支持直接指定表名,或嵌套SQL查询) df = spark.read.jdbc( url=jdbc_url, table="your_target_table", # 示例:"(SELECT id, name FROM users WHERE active = 1) AS active_users" properties=jdbc_properties ) # 验证数据加载结果 df.show(5) # 后续数据处理逻辑...
4. 集成到Airflow任务
将上述逻辑封装为Airflow PythonOperator的调用函数:
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime def load_mysql_to_spark(): # 上述完整代码逻辑 pass with DAG( dag_id="mysql_to_spark_dag", start_date=datetime(2024, 1, 1), schedule_interval="@daily", catchup=False ) as dag: load_task = PythonOperator( task_id="load_mysql_data", python_callable=load_mysql_to_spark )
注意事项
- JDBC驱动依赖:确保Spark环境中存在MySQL JDBC驱动包,可通过
spark.jars配置指定,或提前将驱动包放置在Spark的jars目录下。 - Airflow连接配置:在Airflow UI创建MySQL类型连接时,需正确填写
Host、Port、Login、Password、Schema(数据库名)字段。 - 连接参数调整:根据MySQL服务器配置,修改JDBC URL中的
useSSL、serverTimezone等参数,避免连接报错。
内容的提问来源于stack exchange,提问作者Muhammad Imran Tariq
相关产品推荐
相关产品推荐

