Spark 2.1.1在EMR上读写Redshift时出现类型转换异常
问题背景
你在EMR环境中用Spark 2.1.1读写存储在S3的Redshift数据时,DataFrame能正常创建并识别表字段,但调用count()方法时抛出了java.lang.ClassCastException,具体表现为scala.collection.immutable.List$SerializationProxy无法赋值给org.apache.spark.rdd.RDD中类型为scala.collection.Seq的dependencies_字段。
你的代码示例:
scala> :require /home/hadoop/spark-redshift_2.10-2.0.1.jar Added '/home/hadoop/spark-redshift_2.10-2.0.1.jar' to classpath. scala> :require /home/hadoop/RedshiftJDBC41-1.2.12.1017.jar Added '/home/hadoop/RedshiftJDBC41-1.2.12.1017.jar' to classpath. scala> :require /home/hadoop/spark-avro_2.11-3.2.0.jar Added '/home/hadoop/spark-avro_2.11-3.2.0.jar' to classpath. scala> val read_data = (spark.read | .format("com.databricks.spark.redshift") | .option("url", "jdbc:redshift://redshifthost/schema?user=admin&password=password") | .option("query", "SELECT * FROM schema.table LIMIT 1") | .option("tempdir", tempS3Dir) | .option("forward_spark_s3_credentials",true) | .load()) read_data: org.apache.spark.sql.DataFrame = [aid: int, uid: int ... 3 more fields] scala> read_data.count()
异常栈关键信息:
java.lang.ClassCastException: cannot assign instance of scala.collection.immutable.List$SerializationProxy to field org.apache.spark.rdd.RDD.org$apache$spark$rdd$RDD$$dependencies_ of type scala.collection.Seq in instance of org.apache.spark.rdd.MapPartitionsRDD
...(省略后续重复栈信息)
问题原因
核心是Scala版本不兼容导致的序列化冲突:
- Spark 2.1.1是基于Scala 2.11编译的
- 你加载的
spark-redshift_2.10-2.0.1.jar是Scala 2.10版本的,而spark-avro_2.11-3.2.0.jar是Scala 2.11版本的 - 不同Scala版本的集合类(比如
List)序列化实现逻辑不同,混合使用跨版本jar包时,就会出现这种序列化后的类型转换错误
解决方案
统一Scala版本依赖
替换spark-redshift的jar包为Scala 2.11版本的,比如spark-redshift_2.11-2.0.1.jar(确保该版本兼容Spark 2.1.1)。如果找不到对应版本,可以选择和Spark 2.1.1兼容的spark-redshift版本,比如2.0.0或2.1.0的Scala 2.11编译包。清理冲突依赖
移除所有Scala 2.10版本的jar包,确保添加到classpath的所有依赖jar包都是基于Scala 2.11编译的,彻底避免版本混合。EMR环境优化建议
如果你使用的是EMR集群,可以直接用EMR自带的spark-redshift集成(EMR会预装好与集群Spark版本兼容的依赖),不用手动添加jar包,从根源上避免版本冲突。或者在创建EMR集群时,直接指定兼容的Spark和依赖版本组合。
验证步骤
替换jar包后,重新加载依赖并执行代码:
scala> :require /home/hadoop/spark-redshift_2.11-2.0.1.jar scala> :require /home/hadoop/RedshiftJDBC41-1.2.12.1017.jar scala> :require /home/hadoop/spark-avro_2.11-3.2.0.jar // 重新创建DataFrame并执行count() val read_data = spark.read .format("com.databricks.spark.redshift") .option("url", "jdbc:redshift://redshifthost/schema?user=admin&password=password") .option("query", "SELECT * FROM schema.table LIMIT 1") .option("tempdir", tempS3Dir) .option("forward_spark_s3_credentials",true) .load() read_data.count()
内容的提问来源于stack exchange,提问作者Aravind Krishnakumar

