elasticsearch-pyspark指定字段仍返回全部字段,求问题原因
问题分析与解决方案
我来帮你拆解下这个问题——你遇到的情况其实是elasticsearch-spark连接器(也就是你用的org.elasticsearch.spark.sql格式)的参数处理逻辑和直接调用ES REST API的差异导致的:
为什么Postman正常但PySpark不行?
Postman直接调用ES的_search接口时,ES会完整解析请求体里的所有参数,包括_source的字段过滤规则。但elasticsearch-spark在处理es.query参数时,它只把这个参数里的内容当作查询条件来用,并不会识别和执行查询体里的_source、size这类非查询逻辑的配置——这些规则需要通过连接器专门的配置项来单独设置。
正确的修复方式
你需要把_source的过滤规则从es.query的内容里剥离出来,改用连接器的es.read.source参数来指定要返回的字段:
修改后的PySpark代码
# 只保留查询条件部分,移除原来的_source配置 query_body = "{\"query\":{\"ids\":{\"values\":[\"8\",\"9\",\"10\"]}}}" df = self.__sql_context.read.format("org.elasticsearch.spark.sql") \ .option("es.nodes", "localhost") \ .option("es.port", "9200") \ .option("es.query", query_body) \ .option("es.resource", "test-index4/school") \ .option("es.read.metadata", "true") \ .option("es.read.metadata.version", "true") \ .option("es.read.field.as.array.include", "true") \ # 新增这个参数,指定只返回name字段 .option("es.read.source", "name") \ .load()
额外说明
- 如果需要返回多个字段,用逗号分隔即可:
.option("es.read.source", "name,age,grade") - 如果需要排除某些字段,可以用
es.read.source.exclude参数,比如:.option("es.read.source.exclude", "create_time,update_time") - 注意检查你的elasticsearch-spark依赖版本是否和ES 6.2匹配(应该用6.2.x版本),版本不兼容也可能导致这类配置不生效的问题。
这样调整后,连接器会在查询ES时自动带上_source的过滤规则,和你Postman里的请求效果完全一致,只会返回你指定的name字段以及开启的元数据字段。
内容的提问来源于stack exchange,提问作者amol_shaligram
相关产品推荐
相关产品推荐

