Flink 1.13.6 ES7 Sink写入ES8如何配置setApiCompatibilityMode
问题描述
- Flink 1.13.6 版本仅内置 Elasticsearch 7 版本的 Sink 连接器,流处理作业写入 Elasticsearch 8 集群时触发解析错误
- 可选解决路径有两种:
- 路径一:删除现有集群,重新搭建 Elasticsearch 7 集群
- 路径二:为 Elasticsearch Sink 启用客户端兼容模式
- 路径一不可行:现有集群存储大量重要数据,恢复耗时极长,且 Elasticsearch 不支持高版本快照向低版本恢复,因此只能采用启用兼容模式的方案
- 已知 ES 官方提供了 REST 客户端的跨大版本兼容能力,但不清楚代码中具体配置位置,原有实现代码如下:
import java.io.Serializable; import java.util.Date; import java.util.HashMap; import java.util.List; import java.util.Map; import org.apache.flink.api.common.functions.RuntimeContext; import org.apache.flink.api.java.tuple.Tuple4; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonProcessingException; import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper; import org.apache.flink.streaming.connectors.elasticsearch.ActionRequestFailureHandler; import org.apache.flink.streaming.connectors.elasticsearch.ElasticsearchSinkFunction; import org.apache.flink.streaming.connectors.elasticsearch.RequestIndexer; import org.apache.flink.streaming.connectors.elasticsearch7.ElasticsearchSink; import org.apache.flink.streaming.connectors.elasticsearch7.RestClientFactory; import org.apache.flink.util.ExceptionUtils; import org.apache.http.HttpHost; import org.apache.http.auth.AuthScope; import org.apache.http.auth.UsernamePasswordCredentials; import org.apache.http.client.CredentialsProvider; import org.apache.http.impl.client.BasicCredentialsProvider; import org.apache.http.impl.nio.client.HttpAsyncClientBuilder; import org.elasticsearch.ElasticsearchParseException; import org.elasticsearch.action.ActionRequest; import org.elasticsearch.action.index.IndexRequest; import org.elasticsearch.client.Requests; import org.elasticsearch.client.RestClientBuilder; import org.elasticsearch.common.util.concurrent.EsRejectedExecutionException; public class ElasticMemberSink implements Serializable { private static final long serialVersionUID = 1L; private transient ElasticsearchSink.Builder<Tuple4<String, Integer, Integer, Integer>> memberEsSinkBuilder; public ElasticMemberSink(List<HttpHost> httpHosts, String elasticPassword) { // 初始化ES Sink构造器 memberEsSinkBuilder = new ElasticsearchSink.Builder<>(httpHosts, new ElasticsearchSinkFunction<Tuple4<String, Integer, Integer, Integer>>() { public IndexRequest createIndexRequest(Tuple4<String, Integer, Integer, Integer> memberSummaryTuple) throws JsonProcessingException { Date date = new Date(); Map<String, Object> json = new HashMap<>(); json.put("serverId", memberSummaryTuple.f0); json.put("date", String.valueOf(date.getTime())); json.put("numLeft", memberSummaryTuple.f1); json.put("numJoined", memberSummaryTuple.f2); json.put("memberCount", memberSummaryTuple.f3); return Requests.indexRequest().index("prod-members").type("_doc").source(json); } @Override public void process(Tuple4<String, Integer, Integer, Integer> memberSummaryTuple, RuntimeContext ctx, RequestIndexer indexer) { try { indexer.add(createIndexRequest(memberSummaryTuple)); } catch (JsonProcessingException e) { e.printStackTrace(); } } }); // 配置带认证的REST客户端 if(elasticPassword != null) { memberEsSinkBuilder.setRestClientFactory(restClientBuilder -> { restClientBuilder.setHttpClientConfigCallback(new RestClientBuilder.HttpClientConfigCallback() { @Override public HttpAsyncClientBuilder customizeHttpClient(HttpAsyncClientBuilder httpClientBuilder) { // ES账号密码配置 CredentialsProvider credentialsProvider = new BasicCredentialsProvider(); credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials("elastic", elasticPassword)); return httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider); } }); }); } // 配置批量写入触发条件:每1条事件就刷写 memberEsSinkBuilder.setBulkFlushMaxActions(1); // 配置写入失败处理逻辑 memberEsSinkBuilder.setFailureHandler(new ActionRequestFailureHandler() { @Override public void onFailure(ActionRequest action, Throwable failure, int restStatusCode, RequestIndexer indexer) throws Throwable { if (ExceptionUtils.findThrowable(failure, EsRejectedExecutionException.class).isPresent()) { // 队列满,重新加入索引队列重试 indexer.add(action); } else if (ExceptionUtils.findThrowable(failure, ElasticsearchParseException.class).isPresent()) { // 文档格式错误,直接丢弃不触发作业失败 } else { // 其他错误直接抛出,触发作业失败 throw failure; } } }); } public ElasticsearchSink.Builder<Tuple4<String, Integer, Integer, Integer>> getSinkBuilder() { return memberEsSinkBuilder; } }
实现方法
兼容模式的配置位置就在自定义RestClientFactory的代码块中,核心逻辑是给REST客户端添加统一的兼容请求头,让Elasticsearch 8服务端收到7版本客户端的请求时,按照7版本的协议规则解析请求、返回响应,从根源上避免版本不匹配导致的解析错误。
具体修改步骤:
- 导入请求头相关类:
org.apache.http.Header、org.apache.http.message.BasicHeader - 在
setRestClientFactory的逻辑中,调用restClientBuilder.setDefaultHeaders()方法,传入两个固定的兼容请求头 - 原有业务写入逻辑、认证逻辑、失败处理逻辑都不需要改动
修改后的setRestClientFactory代码片段如下:
// 配置带认证+兼容模式的REST客户端 if(elasticPassword != null) { memberEsSinkBuilder.setRestClientFactory(restClientBuilder -> { // 配置ES8兼容请求头 Header[] defaultHeaders = new Header[]{ new BasicHeader("Accept", "application/vnd.elasticsearch+json;compatible-with=7"), new BasicHeader("Content-Type", "application/vnd.elasticsearch+json;compatible-with=7") }; restClientBuilder.setDefaultHeaders(defaultHeaders); restClientBuilder.setHttpClientConfigCallback(httpClientBuilder -> { // ES账号密码配置 CredentialsProvider credentialsProvider = new BasicCredentialsProvider(); credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials("elastic", elasticPassword)); return httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider); }); }); }
注意:原有代码中写入请求使用的
.type("_doc")写法在Elasticsearch 8中完全兼容,不需要额外删除修改。配置完成后重启作业即可正常写入ES8集群,不需要改动现有集群的任何配置、也不需要迁移数据。
内容的提问来源于stack exchange,提问作者santiar
相关产品推荐
相关产品推荐

