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

Airflow部署EMR Spark Job无法找到MySQL Connector问题排查

EMR集群运行Spark任务时找不到MySQL驱动包的排查与解决

报错信息:

java.io.FileNotFoundException: File file:/mnt/var/lib/hadoop/steps/[STEP-ID]/mysql-connector-j-8.0.33.jar does not exist

以下是脱敏后的配置代码:

一、EMR集群配置(已配置Bootstrap动作)

# Spark Configurations
JOB_FLOW_OVERRIDES = {
    "Name": "EMR spark test",
    "ReleaseLabel": "emr-7.0.0",
    "Applications": [{"Name": "Hadoop"}, {"Name": "Spark"}],
    "Configurations": [
        {
            "Classification": "spark-env",
            "Configurations": [
                {
                    "Classification": "export",
                    "Properties": {"PYSPARK_PYTHON": "/usr/bin/python3"},
                }
            ],
        }
    ],
    "Instances": {
        "InstanceGroups": [
            {
                "Name": "Master node",
                "Market": "SPOT",
                "InstanceRole": "MASTER",
                "InstanceType": "m5.xlarge",
                "InstanceCount": 1,
            },
            {
                "Name": "Core - 2",
                "Market": "SPOT",
                "InstanceRole": "CORE",
                "InstanceType": "m5.xlarge",
                "InstanceCount": 2,
            },
            
        ],
        "KeepJobFlowAliveWhenNoSteps": True,
        "TerminationProtected": False,
        "Ec2SubnetId": "{{ var.value.ec2_subnet_id }}",
    },
    "BootstrapActions": [
        {
            "Name": "import custom Jars",
            "ScriptBootstrapAction": {
                "Path": "s3://PATH/copy_jars.sh",
                "Args": []
            }
        }
    ],
    "JobFlowRole": "EMR_EC2_DefaultRole",
    "ServiceRole": "EMR_DefaultRole_V2",
    "LogUri": "s3://PATH/",
    
}

二、Spark任务步骤配置

# Steps to run
SPARK_STEPS = [
    {
        "Name": "workflow_extraction",
        "ActionOnFailure": "CANCEL_AND_WAIT",
        "HadoopJarStep": {
            "Jar": "command-runner.jar",
            "Args": [
                "spark-submit",
                "--master",
                "yarn",
                "--deploy-mode",
                "client",
                "--jars",
                "mysql-connector-j-8.0.33.jar",
                "--driver-class-path",
                "mysql-connector-j-8.0.33.jar",
                "--conf",
                "spark.executor.extraClassPath=mysql-connector-j-8.0.33.jar",
                "s3://PATH/spark_test.py",
                
                
            ],
        },
    }
]

三、PySpark代码

from pyspark.sql import SparkSession

bucket = "s3://PATH/"

spark = SparkSession.builder \
    .appName("SparkWorkflow") \
    .getOrCreate()

url = "jdbc:mysql://host:3306/database"

mysql_properties = {
    "user": "user",
    "password": "password",
    "driver": "com.mysql.jdbc.Driver"
}

table_name = "table"

df = spark.read.jdbc(url, table_name, properties=mysql_properties)

df.write.parquet(bucket,mode="overwrite")

四、Airflow DAG代码

from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.dummy import DummyOperator
from airflow.providers.amazon.aws.operators.emr import EmrCreateJobFlowOperator, EmrAddStepsOperator, EmrTerminateJobFlowOperator
from airflow.providers.amazon.aws.sensors.emr import EmrStepSensor

# Set default arguments
default_args = {
    "owner": "airflow",
    "start_date": datetime(2022, 3, 5),
    "email": ["airflow@airflow.com"],
    "email_on_failure": False,
    "email_on_retry": False,
    "retries": 1,
    "retry_delay": timedelta(minutes=5),
}

# DAG
with DAG(
    "emr_and_airflow_integration",
    default_args=default_args,
    schedule_interval="0 1 * * *",
    max_active_runs=1,
    catchup=False
) as dag:

    start_dag = DummyOperator(task_id="start_dag")

    create_emr_cluster = EmrCreateJobFlowOperator(
        task_id="create_emr_cluster",
        job_flow_overrides=JOB_FLOW_OVERRIDES,
        aws_conn_id='aws_default',
        emr_conn_id='emr_default',
        region_name='us-east-1'
    )

    add_steps = EmrAddStepsOperator(
        task_id="add_steps",
        job_flow_id="{{ task_instance.xcom_pull(task_ids='create_emr_cluster', key='return_value') }}",
        aws_conn_id="aws_default",
        steps=SPARK_STEPS
    )
    last_step = len(SPARK_STEPS) - 1

    check_execution_steps = EmrStepSensor(
        task_id="check_execution_steps",
        job_flow_id="{{ task_instance.xcom_pull('create_emr_cluster', key='return_value') }}",
        step_id="{{ task_instance.xcom_pull(task_ids='add_steps', key='return_value')["
        + str(last_step)
        + "] }}",
        aws_conn_id="aws_default",
    )

    terminate_emr_cluster = EmrTerminateJobFlowOperator(
        task_id="terminate_emr_cluster",
        job_flow_id="{{ task_instance.xcom_pull(task_ids='create_emr_cluster', key='return_value') }}",
        aws_conn_id="aws_default",
    )

    end_dag = DummyOperator(task_id="end_dag")

    start_dag >> create_emr_cluster >> add_steps
    add_steps >> check_execution_steps >> terminate_emr_cluster >> end_dag

错误原因与修复方案

1. Bootstrap脚本未将驱动包放到Spark可访问路径

你配置了Bootstrap动作拉取驱动包,但copy_jars.sh脚本大概率没把jar包放到Spark任务能识别的全局路径,比如/usr/lib/spark/jars/,导致Spark在步骤工作目录找不到文件。

修复:修改copy_jars.sh脚本,将S3上的驱动包复制到集群所有节点的Spark jars目录:

#!/bin/bash
aws s3 cp s3://PATH/mysql-connector-j-8.0.33.jar /usr/lib/spark/jars/
chmod 755 /usr/lib/spark/jars/mysql-connector-j-8.0.33.jar

同时确保EMR_EC2_DefaultRole有访问该S3路径的权限。

2. Spark Submit参数未指定驱动包绝对路径

当前--jars、--driver-class-path等参数只写了jar文件名,Spark默认在当前步骤工作目录查找,但Bootstrap脚本未将jar放到该目录。

修复:如果已通过Bootstrap将jar放到/usr/lib/spark/jars/,可修改Spark Submit参数(该目录属于Spark默认classpath,--jars参数可省略):

"Args": [
    "spark-submit",
    "--master",
    "yarn",
    "--deploy-mode",
    "client",
    "--driver-class-path",
    "/usr/lib/spark/jars/mysql-connector-j-8.0.33.jar",
    "--conf",
    "spark.executor.extraClassPath=/usr/lib/spark/jars/mysql-connector-j-8.0.33.jar",
    "s3://PATH/spark_test.py",
]

3. MySQL驱动类名不匹配(针对8.0版本)

你使用的是mysql-connector-j-8.0.33.jar,对应的驱动类应为com.mysql.cj.jdbc.Driver,旧版com.mysql.jdbc.Driver在8.0中已废弃,可能导致加载失败。

修复:修改PySpark代码中的驱动类:

mysql_properties = {
    "user": "user",
    "password": "password",
    "driver": "com.mysql.cj.jdbc.Driver"
}

4. 备选方案:直接引用S3上的驱动包

如果不想用Bootstrap脚本,可直接在--jars参数中指定S3路径,Spark会自动下载到集群节点:

"Args": [
    "spark-submit",
    "--master",
    "yarn",
    "--deploy-mode",
    "client",
    "--jars",
    "s3://PATH/mysql-connector-j-8.0.33.jar",
    "--driver-class-path",
    "mysql-connector-j-8.0.33.jar",
    "--conf",
    "spark.executor.extraClassPath=mysql-connector-j-8.0.33.jar",
    "s3://PATH/spark_test.py",
]

需确保EMR集群角色有访问该S3路径的权限。


内容的提问来源于stack exchange,提问作者Fábio Elias Reis Ritter

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 21:06:01