You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 02:00:38