使用QueryElasticsearchHttp从Elasticsearch8.6.1同步数据到NiFi1.19.1遇类型错误
解决Apache NiFi QueryElasticsearchHttp连接Elasticsearch 8.x的类型问题
问题根源
Elasticsearch 7.x后废弃了索引类型(_type),8.x完全移除该概念,但NiFi 1.19.1的QueryElasticsearchHttp处理器源码仍强制要求配置Type属性,导致两种矛盾错误:
- 不配置Type:触发NullPointerException
- 配置任意Type值:ES返回400 Bad Request(8.x不支持
_type参数)
可行解决方案
方案1:调整处理器配置适配无类型场景
在QueryElasticsearchHttp的Type属性中填入空字符串(不能留空不填),同时修改Request Path属性:
- 将默认的
/${index}/${type}/_search改为/${index}/_search,直接移除路径中的类型占位符。
这样既满足处理器对Type属性非空的校验,又符合ES 8.x的请求格式,避免400错误。
方案2:升级NiFi版本
NiFi 1.20.0及后续版本已适配Elasticsearch 8.x的无类型特性,QueryElasticsearchHttp不再强制要求Type属性,默认请求路径也已调整为不含_type的格式,直接升级即可彻底解决问题。
方案3:自定义脚本绕过处理器限制(临时 workaround)
如果暂时无法升级NiFi,可使用ExecuteScript处理器编写Groovy脚本直接调用ES REST API,示例核心逻辑:
import org.apache.http.client.methods.HttpGet import org.apache.http.impl.client.CloseableHttpClient import org.apache.http.impl.client.HttpClients import org.apache.http.util.EntityUtils def client = HttpClients.createDefault() def esUrl = "http://your-es-cluster:9200/your-target-index/_search?q=*" def request = new HttpGet(esUrl) request.addHeader("Content-Type", "application/json") // 若ES开启安全验证,添加认证头 // request.addHeader("Authorization", "Bearer your-access-token") def response = client.execute(request) def responseBody = EntityUtils.toString(response.getEntity()) def flowFile = session.create() flowFile = session.write(flowFile, { outputStream -> outputStream.write(responseBody.getBytes("UTF-8")) } as OutputStreamCallback) session.transfer(flowFile, REL_SUCCESS) client.close()
验证要点
- 启动处理器后,查看NiFi日志确认无报错
- 检查输出FlowFile是否包含正确的ES查询结果
- 查看ES访问日志,确认请求路径为
/_search而非带类型的格式
内容的提问来源于stack exchange,提问作者Malinda Peiris
相关产品推荐
相关产品推荐

