Spark读取Elasticsearch时Struct字段找不到报错如何解决
问题原因分析
- 版本兼容性问题:你使用的Spark 3.1.2对应Scala版本为2.12,但当前引入的elasticsearch-spark依赖包为
elasticsearch-spark-20_2.11-7.9.0.jar,其中_2.11代表该Jar包基于Scala 2.11编译,不同大版本的Scala编译产物互不兼容;同时spark-20代表该包适配Spark 2.x版本,和你使用的Spark 3.x版本也存在适配冲突,这是核心诱因之一。 - 字段映射不匹配:错误提示明确说明
Field 'about' not found; typically this occurs with arrays which are not mapped as single value,说明store_products索引中部分文档的about字段实际为数组类型,但ES连接器自动推导Schema时将其识别为单值Struct类型,读取到数组类型的文档时就会触发解析失败。
解决方案
- 首先替换适配版本的ES连接器Jar包,使用
elasticsearch-spark-30_2.12-7.9.0.jar,满足两个适配要求:- 后缀
_2.12匹配Spark使用的Scala 2.12版本 - 前缀
spark-30适配Spark 3.x版本
- 后缀
- 调整Spark读取ES的配置,根据字段实际属性增加对应参数:
- 若确认
about字段为数组类型,增加配置option("es.read.field.as.array.include", "about"),明确告知连接器该字段为数组类型,避免Schema推导错误 - 若确认
about字段应为单值Struct,仅部分文档存在格式异常,可增加配置option("es.spark.dataframe.null.on.array", "true"),遇到数组类型时直接返回null避免任务崩溃
- 若确认
- 修复后的代码示例:
val esURL = "ES.com" val reader = spark.read .format("org.elasticsearch.spark.sql") .option("es.nodes.wan.only","true") .option("es.port","443") .option("es.net.ssl","true") .option("es.nodes", esURL) .option("es.net.http.auth.user","admin") .option("es.net.http.auth.pass","pass") // 新增配置,指定about为数组类型,可根据实际需求调整 .option("es.read.field.as.array.include", "about") val df = reader.load("store_products") df.printSchema() // 若about为数组类型,需先使用explode展开后再提取内部字段 import org.apache.spark.sql.functions.explode df.select(explode($"about").alias("about_item")).select("about_item.*").show()
内容的提问来源于stack exchange,提问作者TheDataGuy
相关产品推荐
相关产品推荐

