You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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版本的协议规则解析请求、返回响应,从根源上避免版本不匹配导致的解析错误。

具体修改步骤:

  1. 导入请求头相关类:org.apache.http.Header、org.apache.http.message.BasicHeader
  2. 在setRestClientFactory的逻辑中,调用restClientBuilder.setDefaultHeaders()方法,传入两个固定的兼容请求头
  3. 原有业务写入逻辑、认证逻辑、失败处理逻辑都不需要改动

修改后的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 12:51:31