使用elasticsearch-hadoop连接Spark时遇scala.Product$class类未找到异常
Spark读取Elasticsearch触发scala.Product$class缺失异常问题
问题场景
在AWS EMR上运行spark-submit任务读取Elasticsearch节点数据时,触发java.lang.NoClassDefFoundError: scala/Product$class异常。环境配置:
- EMR版本:3.3.2-amzn-0
- Scala版本:2.12.15
- Java版本:1.8.0_382
代码示例
Python代码
es_config = { "es.nodes": url_to_my_node, "es.port": "9200", "es.resource": "my_elasticsearch_index/_doc", "es.query": "?q=id:park_rocky-mountain", "es.read.metadata": "true", "es.nodes.wan.only": "true", } df = spark.read.format("org.elasticsearch.spark.sql") \ .options(**es_config) \ .load()
Scala代码
val esReadOptions = Map( "es.nodes" -> "url_to_my_node", "es.port" -> "9200" ) val df = spark.read.format("org.elasticsearch.spark.sql") .options(esReadOptions) .load("url_to_my_node")
异常信息
java.lang.NoClassDefFoundError: scala/Product$class at org.elasticsearch.spark.sql.ElasticsearchRelation.<init>(DefaultSource.scala:221) at org.elasticsearch.spark.sql.DefaultSource.createRelation(DefaultSource.scala:97) at org.apache.spark.sql.execution.datasources.DataSource.resolveRelation(DataSource.scala:350) at org.apache.spark.sql.DataFrameReader.loadV1Source(DataFrameReader.scala:228) at org.apache.spark.sql.DataFrameReader.$anonfun$load$2(DataFrameReader.scala:210) at scala.Option.getOrElse(Option.scala:189) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:210) at org.apache.spark.sql.DataFrameReader.load(DataFrameReader.scala:171) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182) at py4j.ClientServerConnection.run(ClientServerConnection.java:106) at java.lang.Thread.run(Thread.java:750) Caused by: java.lang.ClassNotFoundException: scala.Product$class at java.net.URLClassLoader.findClass(URLClassLoader.java:387) at java.lang.ClassLoader.loadClass(ClassLoader.java:418) at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:352) at java.lang.ClassLoader.loadClass(ClassLoader.java:351)
已执行的排查动作
- 确认Spark上下文中存在
elasticsearch-hadoop-7.17.6.jar,该版本本应兼容Scala 2.12与Spark 3.x - 尝试更换多个版本的JAR包:
elasticsearch-hadoop-8.9.0.jar、7.17.6.jar、opensearch-hadoop-1.0.1.jar - 连接不同的Elasticsearch实例:8.9.0版本沙箱、兼容ES7.12的OpenSearch实例,均已通过REST API验证可正常访问
- 在PySpark、Scala spark-shell中测试,AWS EMR及本地Linux Mint环境下均出现相同异常
疑问
scala.Product$class是Scala 2.3及以上版本的标准库内容,当前使用Scala 2.12版本却触发该类缺失异常,怀疑错误信息可能隐藏了其他底层问题。
内容的提问来源于stack exchange,提问作者whenthoughsmatch
相关产品推荐
相关产品推荐

