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

添加Elasticsearch Appender至org.apache.http日志后应用挂起

自定义Elasticsearch Logback Appender导致应用挂起的问题分析与解决方案

问题场景

我们实现了自定义Logback Elasticsearch Appender,代码如下:

//
// Source code recreated from a .class file by IntelliJ IDEA
// (powered by FernFlower decompiler)
//

import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.core.AppenderBase;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.io.IOException;
import java.net.InetAddress;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.ZoneOffset;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import org.apache.commons.codec.binary.Base64;
import org.apache.http.HttpHost;
import org.apache.http.impl.nio.reactor.IOReactorConfig;
import org.elasticsearch.ElasticsearchStatusException;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestClient;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.client.indices.CreateIndexRequest;
import org.elasticsearch.common.xcontent.XContentType;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

public class ElasticSearchAppender extends AppenderBase<ILoggingEvent> {
    private static final Logger log = LoggerFactory.getLogger(ElasticSearchAppender.class);
    private static ObjectMapper mapper = new ObjectMapper();
    private static RestHighLevelClient elastic;
    private ExecutorService executorService = Executors.newFixedThreadPool(10);
    private static RequestOptions COMMON_OPTIONS;
    private String hostname;
    private Integer port;
    private String index;
    private String application;
    private String login;
    private String password;

    public ElasticSearchAppender() {
    }

    public String getLogin() {
        return this.login;
    }

    public void setLogin(String login) {
        this.login = login;
    }

    public String getPassword() {
        return this.password;
    }

    public void setPassword(String password) {
        this.password = password;
    }

    public String getApplication() {
        return this.application;
    }

    public void setApplication(String application) {
        this.application = application;
    }

    public String getHostname() {
        return this.hostname;
    }

    public void setHostname(String hostname) {
        this.hostname = hostname;
    }

    public Integer getPort() {
        return this.port;
    }

    public void setPort(Integer port) {
        this.port = port;
    }

    public String getIndex() {
        return this.index;
    }

    public void setIndex(String index) {
        this.index = index;
    }

    protected void append(ILoggingEvent event) {
        if (elastic == null) {
            this.createElasticConnection(this.hostname, this.port);

            try {
                CreateIndexRequest createIndexRequest = new CreateIndexRequest(this.index);
                createIndexRequest.source(mapper.writeValueAsString(new Message(InetAddress.getLocalHost().getHostName(), LocalDateTime.ofInstant(Instant.ofEpochMilli(event.getTimeStamp()), ZoneOffset.UTC).toString(), event.getLevel().toString(), event.getLoggerName(), event.getMDCPropertyMap(), event.getFormattedMessage(), this.application)), XContentType.JSON);
                elastic.indices().create(createIndexRequest, COMMON_OPTIONS);
            } catch (ElasticsearchStatusException var3) {
                log.error(var3.getMessage(), var3);
            } catch (IOException var4) {
            }
        } else {
            this.executorService.submit(() -> {
                try {
                    IndexRequest indexRequest = new IndexRequest(this.index);
                    indexRequest.source(mapper.writeValueAsString(new Message(InetAddress.getLocalHost().getHostName(), LocalDateTime.ofInstant(Instant.ofEpochMilli(event.getTimeStamp()), ZoneOffset.UTC).toString(), event.getLevel().toString(), event.getLoggerName(), event.getMDCPropertyMap(), event.getFormattedMessage(), this.application)), XContentType.JSON);
                    elastic.index(indexRequest, COMMON_OPTIONS);
                } catch (IOException var3) {
                }

            });
        }

    }

    private void createElasticConnection(String hostname, int port) {
        elastic = new RestHighLevelClient(RestClient.builder(new HttpHost[]{new HttpHost(hostname, port, "http")}).setRequestConfigCallback((requestConfigBuilder) -> {
            return requestConfigBuilder.setConnectTimeout(5000).setSocketTimeout(60000);
        }).setHttpClientConfigCallback((httpClientBuilder) -> {
            return httpClientBuilder.setDefaultIOReactorConfig(IOReactorConfig.custom().setIoThreadCount(200).build());
        }));
        RequestOptions.Builder builder = RequestOptions.DEFAULT.toBuilder();
        builder.addHeader("Authorization", "Basic " + Base64.encodeBase64String((this.login + ":" + this.password).getBytes()));
        COMMON_OPTIONS = builder.build();
    }

    private class Message {
        private String hostname;
        private String time;
        private String level;
        private String loggerName;
        private Map<String, String> mdcProperty;
        private String msg;
        private String applicationName;

        public String getHostname() {
            return this.hostname;
        }

        public String getTime() {
            return this.time;
        }

        public String getLevel() {
            return this.level;
        }

        public String getLoggerName() {
            return this.loggerName;
        }

        public Map<String, String> getMdcProperty() {
            return this.mdcProperty;
        }

        public String getMsg() {
            return this.msg;
        }

        public String getApplicationName() {
            return this.applicationName;
        }

        public void setHostname(String hostname) {
            this.hostname = hostname;
        }

        public void setTime(String time) {
            this.time = time;
        }

        public void setLevel(String level) {
            this.level = level;
        }

        public void setLoggerName(String loggerName) {
            this.loggerName = loggerName;
        }

        public void setMdcProperty(Map<String, String> mdcProperty) {
            this.mdcProperty = mdcProperty;
        }

        public void setMsg(String msg) {
            this.msg = msg;
        }

        public void setApplicationName(String applicationName) {
            this.applicationName = applicationName;
        }

        public Message() {
        }

        public Message(String hostname, String time, String level, String loggerName, Map<String, String> mdcProperty, String msg, String applicationName) {
            this.hostname = hostname;
            this.time = time;
            this.level = level;
            this.loggerName = loggerName;
            this.mdcProperty = mdcProperty;
            this.msg = msg;
            this.applicationName = applicationName;
        }
    }
}

Logback配置如下:

<appender name="Elastic" class="***.log.appender.ElasticSearchAppender">
    <hostname>${elkHostname}</hostname>
    <port>${elkPort}</port>
    <index>${elkIndex}</index>
    <application>${applicationName}</application>
    <login>${login}</login>
    <password>${password}</password>
</appender>

启用以下Logger配置后,应用在首次接收到或发起HTTP调用时挂起:

<logger name="org.apache.http" level="debug" additivity="false">
     <appender-ref ref="STDOUT"/>
     <appender-ref ref="Elastic"/>
 </logger>

问题原因

  1. 循环日志依赖引发线程阻塞:Elasticsearch RestHighLevelClient底层依赖Apache HttpClient,当org.apache.http的DEBUG日志触发时,自定义Appender会尝试同步初始化ES连接并创建索引,而ES客户端的HTTP请求又会生成新的org.apache.http日志,形成循环调用。此时处理日志的线程被同步的ES初始化操作占用,无法处理新的日志,最终导致线程阻塞,应用挂起。
  2. 非线程安全的单例初始化:静态变量elastic的初始化未加同步控制,多线程环境下可能出现重复初始化或半初始化状态,加剧阻塞风险。
  3. 空异常处理掩盖问题:代码中存在空的IOException catch块,无法捕获初始化或日志写入过程中的异常,无法定位问题根源。

解决方案

1. 切断循环日志依赖

修改Logback配置,禁止将Elastic Appender关联到org.apache.http或ES客户端相关的Logger:

<logger name="org.apache.http" level="debug" additivity="false">
    <appender-ref ref="STDOUT"/>
    <!-- 移除Elastic Appender引用 -->
</logger>

或者在Appender中添加过滤规则,直接拒绝ES和HttpClient相关日志:

<appender name="Elastic" class="***.log.appender.ElasticSearchAppender">
    <!-- 原有配置 -->
    <filter class="ch.qos.logback.core.filter.EvaluatorFilter">
        <evaluator>
            <expression>loggerName.startsWith("org.apache.http") || loggerName.startsWith("org.elasticsearch")</expression>
        </evaluator>
        <onMatch>DENY</onMatch>
        <onMismatch>ACCEPT</onMismatch>
    </filter>
</appender>

2. 异步初始化ES客户端与索引

重写Appender的start()方法,在Appender启动时异步完成ES连接初始化和索引创建,避免阻塞日志线程:

@Override
public void start() {
    super.start();
    executorService.submit(() -> {
        if (elastic == null) {
            synchronized (ElasticSearchAppender.class) {
                if (elastic == null) {
                    try {
                        createElasticConnection(hostname, port);
                        // 使用预定义的索引映射模板,而非日志Message结构
                        CreateIndexRequest createIndexRequest = new CreateIndexRequest(index);
                        String indexMapping = "{\n" +
                                "  \"mappings\": {\n" +
                                "    \"properties\": {\n" +
                                "      \"hostname\": {\"type\": \"keyword\"},\n" +
                                "      \"time\": {\"type\": \"date\", \"format\": \"yyyy-MM-dd'T'HH:mm:ss\"},\n" +
                                "      \"level\": {\"type\": \"keyword\"},\n" +
                                "      \"loggerName\": {\"type\": \"keyword\"},\n" +
                                "      \"mdcProperty\": {\"type\": \"object\"},\n" +
                                "      \"msg\": {\"type\": \"text\"},\n" +
                                "      \"applicationName\": {\"type\": \"keyword\"}\n" +
                                "    }\n" +
                                "  }\n" +
                                "}";
                        createIndexRequest.source(indexMapping, XContentType.JSON);
                        elastic.indices().create(createIndexRequest, COMMON_OPTIONS);
                    } catch (ElasticsearchStatusException e) {
                        log.error("索引已存在或创建失败", e);
                    } catch (IOException e) {
                        log.error("ES连接或索引创建异常", e);
                    }
                }
            }
        }
    });
}

同时修改append()方法,移除初始化逻辑,直接提交日志任务:

protected void append(ILoggingEvent event) {
    if (elastic == null) {
        log.warn("ES客户端尚未初始化,日志暂无法写入");
        return;
    }
    executorService.submit(() -> {
        try {
            IndexRequest indexRequest = new IndexRequest(index);
            indexRequest.source(mapper.writeValueAsString(new Message(
                    InetAddress.getLocalHost().getHostName(),
                    LocalDateTime.ofInstant(Instant.ofEpochMilli(event.getTimeStamp()), ZoneOffset.UTC).toString(),
                    event.getLevel().toString(),
                    event.getLoggerName(),
                    event.getMDCPropertyMap(),
                    event.getFormattedMessage(),
                    application
            )), XContentType.JSON);
            elastic.index(indexRequest, COMMON_OPTIONS);
        } catch (IOException e) {
            log.error("日志写入ES失败", e);
        }
    });
}

3. 修复异常处理

移除所有空catch块,添加异常日志输出,便于问题排查:

// 替换原空catch块
catch (IOException e) {
    log.error("ES操作异常", e);
}

4. 添加资源清理逻辑

重写stop()方法,关闭线程池和ES客户端,避免资源泄漏:

@Override
public void stop() {
    super.stop();
    executorService.shutdown();
    try {
        if (elastic != null) {
            elastic.close();
        }
    } catch (IOException e) {
        log.error("关闭ES客户端失败", e);
    }
}

内容的提问来源于stack exchange,提问作者Eljah

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 10:07:36