Spark Scala调用Elasticsearch查询遇映射不一致错误求助
问题分析
你遇到的EsHadoopIllegalStateException错误提示“Position for 'xxxx' not found in row”,本质是Spark在读取Elasticsearch返回的结果时,找不到某个指定字段(就是错误里的xxxx),通常由以下几种情况导致:
- 字段映射不匹配:Elasticsearch索引里的字段类型(比如嵌套类型、text/keyword类型)和Spark DataFrame的Schema定义不一致;或者字段名大小写不匹配(ES对字段名大小写敏感,Spark如果写错就会找不到)。
- 部分文档缺失字段:你的ES索引里有些文档没有
xxxx字段,但Spark的Schema里这个字段被定义为必填(非nullable),读取到缺失字段的文档就会报错。 - 版本兼容性问题:你使用的
es-spark依赖版本和Elasticsearch集群版本不匹配,导致字段解析逻辑出现冲突。
结合你的查询语句来看,你用query_string: "*"会返回所有字段,很容易碰到上述某一种情况。
具体解决步骤
1. 先定位问题字段
先把错误里的xxxx字段明确下来——这是排查的关键。你可以先简化查询,比如只返回几个确定存在的字段(比如date和keyword),看看是否还报错;或者用ES的查询语句过滤出肯定包含xxxx字段的文档,测试是否能正常读取,以此确认是不是该字段的问题。
2. 检查字段映射与Spark Schema
- 查看ES索引映射:在Kibana Dev Tools或者ES命令行里执行
GET /你的索引名/_mapping,逐一对比字段的类型和Spark代码里的Schema定义。比如ES里是long类型的date,Spark里要对应LongType;如果是嵌套对象,Spark里要定义成StructType。 - 注意字段名大小写:确保Spark代码里引用的字段名和ES里完全一致,比如ES里是
Keyword,Spark里不能写成keyword。
3. 处理缺失字段的情况
如果是部分文档缺失xxxx字段,可以做以下调整:
- 在Spark读取ES数据时,配置
es.read.field.as.array.include参数,把可能缺失的字段指定为数组类型(允许为空); - 在Spark的Schema定义里,把该字段设为可空,比如:
import org.apache.spark.sql.types._ val customSchema = StructType(Array( StructField("date", LongType, nullable = true), StructField("keyword", StringType, nullable = true), StructField("xxxx", StringType, nullable = true) // 设为可空 ))
- 或者修改ES查询语句,明确指定要返回的字段,避免返回那些可能有问题的字段,比如给查询加上
_source:
{ "_source": ["date", "keyword"], "bool": { "must": [ {"query_string": {"query": "*", "analyze_wildcard": true}}, {"range": {"date": {"gte": "1527577200000", "lte": "1527606000000", "format": "epoch_millis"}}}, {"bool": {"must": [{"term": {"keyword": "name=cat"}}]}} ], "must_not": [] } }
4. 验证版本兼容性
确保你的es-spark依赖版本和Elasticsearch集群的主版本完全一致(比如ES集群是7.17.5,那依赖也要用7.17.5版本),版本不匹配很容易导致各种解析异常。
内容的提问来源于stack exchange,提问作者Bab
相关产品推荐
相关产品推荐

