Airflow调度PySpark ETL任务时触发java.io.InvalidClassException错误求助
我正在构建一套ETL流程:从API获取数据→PySpark分布式转换→上传至MongoDB,并用Airflow实现自动化调度。完成Spark、Airflow组件部署后,Spark任务能被触发但执行报错,错误信息为java.io.InvalidClassException: org.apache.spark.scheduler.Task; local class incompatible,具体表现为序列化版本号不匹配。已确认所有Spark Master和Worker节点的Spark版本一致,且该任务脱离Airflow单独运行时可正常执行。
错误日志
java.io.InvalidClassException: org.apache.spark.scheduler.Task; local class incompatible: stream classdesc serialVersionUID = 8518222255558883333, local class serialVersionUID = 9998887776665554444
at java.io.ObjectStreamClass.initNonProxy(ObjectStreamClass.java:699)
at java.io.ObjectInputStream.readNonProxyDesc(ObjectInputStream.java:2003)
...(省略堆栈其余部分)
Docker Compose关键配置片段
version: '3.8' services: airflow-webserver: image: apache/airflow:2.8.0 environment: - AIRFLOW__CORE__EXECUTOR=CeleryExecutor - SPARK_HOME=/opt/spark - PYSPARK_PYTHON=/usr/local/bin/python volumes: - ./dags:/opt/airflow/dags - ./spark-jars:/opt/spark/jars - ./scripts:/opt/airflow/scripts depends_on: - airflow-scheduler - spark-master spark-master: image: bitnami/spark:3.5.0 environment: - SPARK_MODE=master ports: - "7077:7077" - "8080:8080" spark-worker: image: bitnami/spark:3.5.0 environment: - SPARK_MODE=worker - SPARK_MASTER_URL=spark://spark-master:7077 depends_on: - spark-master
movies.py脚本核心片段
from pyspark.sql import SparkSession from pyspark.sql.functions import col, current_timestamp import requests def fetch_movie_data(): response = requests.get("https://api.example.com/movies") return response.json() def main(): spark = SparkSession.builder \ .appName("MovieETL") \ .master("spark://spark-master:7077") \ .config("spark.mongodb.output.uri", "mongodb://mongo:27017/movies_db.movies") \ .getOrCreate() movie_data = fetch_movie_data() df = spark.createDataFrame(movie_data) transformed_df = df.withColumn("ingestion_time", current_timestamp()) transformed_df.write.format("mongodb").mode("append").save() spark.stop() if __name__ == "__main__": main()
排查与解决方案
1. 对齐Airflow与Spark集群的PySpark版本
Airflow官方镜像自带的PySpark版本可能和你的Spark集群版本(3.5.0)不匹配,导致序列化类的UID不一致:
- 进入Airflow容器,执行
pip show pyspark查看版本,确认是否与集群版本一致。 - 若版本不符,在容器内安装对应版本:
pip install pyspark==3.5.0,或构建自定义Airflow镜像时预先指定安装该版本的PySpark。
2. 切换至Kryo序列化器
默认Java序列化器对版本差异敏感,改用KryoSerializer可提升兼容性:
spark = SparkSession.builder \ .appName("MovieETL") \ .master("spark://spark-master:7077") \ .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \ .config("spark.mongodb.output.uri", "mongodb://mongo:27017/movies_db.movies") \ .getOrCreate()
3. 移至Spark Worker端执行数据获取
当前脚本在Airflow容器内调用requests获取数据,再传递给Spark,可能引发跨环境类序列化冲突。修改为在Spark Worker节点执行数据获取:
def fetch_movie_data_worker(): import requests response = requests.get("https://api.example.com/movies") return response.json() def main(): spark = SparkSession.builder...getOrCreate() # 让Worker节点执行API调用 movie_data = spark.sparkContext.parallelize([1]).map(lambda x: fetch_movie_data_worker()).collect()[0] df = spark.createDataFrame(movie_data) ...
4. 统一Java版本
序列化异常也可能由Java版本差异导致:
- 检查Airflow容器内的Java版本(
java -version),确保与Spark节点的Java版本一致(Spark 3.5.0推荐Java 11)。
5. 使用SparkSubmitOperator提交任务
避免直接在Airflow环境执行PySpark脚本,改用SparkSubmitOperator让任务完全使用Spark集群环境:
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator from airflow.models import DAG from datetime import datetime default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1) } with DAG('movie_etl_dag', default_args=default_args, schedule_interval='@daily') as dag: spark_task = SparkSubmitOperator( task_id='submit_movie_etl', application='/opt/airflow/scripts/movies.py', conn_id='spark_default', conf={ 'spark.mongodb.output.uri': 'mongodb://mongo:27017/movies_db.movies' }, jars='/opt/spark/jars/mongo-spark-connector_2.12-10.2.0.jar' )
同时在Airflow的Connections中配置spark_default,指向Spark Master地址spark://spark-master:7077。
内容的提问来源于stack exchange,提问作者Astro

