Airflow部署EMR Spark Job无法找到MySQL Connector问题排查
报错信息:
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

