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

Airflow调度PySpark ETL任务时触发java.io.InvalidClassException错误求助

问题:Airflow调度PySpark任务时出现序列化版本不匹配错误(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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 23:32:02