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

如何通过OpenTelemetry将NiFi的Logback自定义Appender日志推送到ElasticSearch

Java自定义Logback Appender对接OpenTelemetry实现方案

需求背景

  • 当前使用Apache NiFi,日志框架为Logback,已实现自定义Appender推送日志到MongoDB
  • 新需求要求将完全相同格式的日志通过OpenTelemetry推送到ElasticSearch,和团队现有技术方案对齐
  • 原有MongoDB推送逻辑保留不变

1. 引入依赖

首先在项目pom.xml中添加OpenTelemetry相关依赖,版本请和团队其他应用使用的OTel版本保持一致,避免兼容性问题:

<dependencies>
    <!-- OpenTelemetry核心API -->
    <dependency>
        <groupId>io.opentelemetry</groupId>
        <artifactId>opentelemetry-api</artifactId>
        <version>1.32.0</version>
    </dependency>
    <!-- OpenTelemetry日志SDK -->
    <dependency>
        <groupId>io.opentelemetry</groupId>
        <artifactId>opentelemetry-sdk-logs</artifactId>
        <version>1.32.0</version>
    </dependency>
    <!-- 全局OTel实例工具 -->
    <dependency>
        <groupId>io.opentelemetry</groupId>
        <artifactId>opentelemetry-extension-annotations</artifactId>
        <version>1.32.0</version>
    </dependency>
</dependencies>

2. 完整Appender实现代码

在原有MongoDB Appender的基础上新增OpenTelemetry上报逻辑,原有业务逻辑完全复用:

import ch.qos.logback.classic.Level;
import ch.qos.logback.classic.spi.ILoggingEvent;
import ch.qos.logback.classic.spi.ThrowableProxyUtil;
import ch.qos.logback.core.AppenderBase;
import com.mongodb.client.MongoClient;
import com.mongodb.client.MongoDatabase;
import com.mongodb.client.model.InsertOneResult;
import io.opentelemetry.api.GlobalOpenTelemetry;
import io.opentelemetry.api.logs.Logger;
import io.opentelemetry.api.logs.LoggerProvider;
import io.opentelemetry.api.logs.Severity;
import io.opentelemetry.api.common.Attributes;
import org.bson.Document;
import org.reactivestreams.Publisher;
import org.reactivestreams.Subscriber;
import org.reactivestreams.Subscription;
import java.util.Date;
import java.util.concurrent.TimeUnit;

public class NiFiCombinedLogAppender extends AppenderBase<ILoggingEvent> {
    // 原有MongoDB相关属性
    private String server;
    private int port;
    private String databaseName;
    private String userName;
    private String password;
    private String collectionName;
    private String hostIp;
    private MongoClient client;
    private MongoDatabase database;

    // 新增OpenTelemetry相关属性
    private Logger otelLogger;

    @Override
    public void start() {
        super.start();
        System.out.println("Initialising Combined LogAppender ~~~~~~~~~~~~~~~~");
        try{
            // 原有MongoDB连接初始化
            createConnection();
            // 新增OpenTelemetry日志实例初始化,直接使用团队统一配置的全局OTel实例
            LoggerProvider loggerProvider = GlobalOpenTelemetry.get().getLogsBridge();
            otelLogger = loggerProvider.loggerBuilder("nifi-log-appender")
                    .setInstrumentationVersion("1.0.0")
                    .build();
        }catch (Exception e){
            System.err.printf("Failed to initialise appender, Mongo server [%s], port [%s]%n", this.server, this.port);
        }
    }

    private MongoClient createConnection () throws Exception {
        if (client == null)
        {
            client = createMongoClient(this.server, this.port, this.databaseName, this.userName, this.password);
        }

        if (this.databaseName != null && (!"".equals(this.databaseName))) {
            database = client.getDatabase(this.databaseName);
        } else {
            System.err.println("Mongo database name is required.");
        }
        return client;
    }


    @Override
    protected void append(ILoggingEvent iLoggingEvent) {
        if (database == null || otelLogger == null)
            return;

        String logMessage = iLoggingEvent.getMessage() == null ? "" : iLoggingEvent.getMessage();
        String[] msgParts = parseLogMessage(logMessage);
        
        // 原有MongoDB日志插入逻辑完全保留
        Document doc = new Document()
                .append("Timestamp", new Date(iLoggingEvent.getTimeStamp()))
                .append("Ip", this.hostIp)
                .append("Server", "NiFi")
                .append("Instance", "")
                .append("Url", "")
                .append("TTId", msgParts[1])
                .append("LTId", msgParts[0])
                .append("LUId", msgParts[2])
                .append("SId", msgParts[4])
                .append("RId", msgParts[3])
                .append("Level", iLoggingEvent.getLevel().levelStr)
                .append("Logger", iLoggingEvent.getLoggerName())
                .append("Thread", iLoggingEvent.getThreadName())
                .append("Message", msgParts[5])
                .append("Exception", iLoggingEvent.getThrowableProxy() != null ? ThrowableProxyUtil.asString(iLoggingEvent.getThrowableProxy()) : null);
        try
        {
            Publisher<InsertOneResult> publisher = database.getCollection(collectionName).insertOne(doc);
            publisher.subscribe(new Subscriber<InsertOneResult>() {
                @Override
                public void onSubscribe(final Subscription s) {
                    s.request(1);
                }

                @Override
                public void onNext(final InsertOneResult result) {}

                @Override
                public void onError(final Throwable t) {
                    System.err.println("Failed to insert Nifi log to mongodb : "+this.toString());
                    t.printStackTrace();
                }

                @Override
                public void onComplete() {}
            });
        }
        catch (Exception e)
        {
            System.err.println("Encountered exception while logging Nifi log to MongoDB : "+this.toString());
            e.printStackTrace();
        }

        // 新增:OpenTelemetry日志上报逻辑,字段和MongoDB完全一致
        try {
            var logRecordBuilder = otelLogger.logRecordBuilder()
                    .setTimestamp(iLoggingEvent.getTimeStamp(), TimeUnit.MILLISECONDS)
                    .setSeverity(convertToOtelSeverity(iLoggingEvent.getLevel()))
                    .setSeverityText(iLoggingEvent.getLevel().levelStr)
                    .setBody(msgParts[5])
                    .setAttribute("Ip", this.hostIp)
                    .setAttribute("Server", "NiFi")
                    .setAttribute("Instance", "")
                    .setAttribute("Url", "")
                    .setAttribute("TTId", msgParts[1])
                    .setAttribute("LTId", msgParts[0])
                    .setAttribute("LUId", msgParts[2])
                    .setAttribute("SId", msgParts[4])
                    .setAttribute("RId", msgParts[3])
                    .setAttribute("Logger", iLoggingEvent.getLoggerName())
                    .setAttribute("Thread", iLoggingEvent.getThreadName());

            if (iLoggingEvent.getThrowableProxy() != null) {
                logRecordBuilder.setAttribute("Exception", ThrowableProxyUtil.asString(iLoggingEvent.getThrowableProxy()));
            }
            logRecordBuilder.emit();
        } catch (Exception e) {
            System.err.println("Encountered exception while sending log to OpenTelemetry : "+this.toString());
            e.printStackTrace();
        }
    }

    // Logback日志级别转OTel日志级别工具方法
    private Severity convertToOtelSeverity(Level level) {
        if (level == Level.ERROR) return Severity.ERROR;
        if (level == Level.WARN) return Severity.WARN;
        if (level == Level.INFO) return Severity.INFO;
        if (level == Level.DEBUG) return Severity.DEBUG;
        if (level == Level.TRACE) return Severity.TRACE;
        return Severity.UNDEFINED_SEVERITY_NUMBER;
    }

    /* 原有日志解析方法完全复用
    * Log message format expected : "~(<LoggedinTenantId>, <TargetTenantId>, <Userid>) ~ <RequestId> ~ <SessionId> ~ <DetailedMessage>"
    * */
    public static String[] parseLogMessage(String logMessage){
        String[] msgParts = new String[6];
        // 原有解析逻辑不变
        return msgParts;
    }

    // 原有MongoClient创建方法完全复用
    private synchronized static MongoClient createMongoClient(String server, int port, String databaseName, String userName, String password)
    {
        // 原有创建逻辑不变
        return MongoClients.create(settings);
    }

    // 原有getter、setter方法不变
    public void setServer(String server) { this.server = server; }
    public void setPort(int port) { this.port = port; }
    public void setDatabaseName(String databaseName) { this.databaseName = databaseName; }
    public void setUserName(String userName) { this.userName = userName; }
    public void setPassword(String password) { this.password = password; }
    public void setCollectionName(String collectionName) { this.collectionName = collectionName; }
    public void setHostIp(String hostIp) { this.hostIp = hostIp; }
}

3. Logback配置示例

<configuration scan="true" scanPeriod="30 seconds">
    <appender name="CombinedLogAppender" class="com.yourpackage.NiFiCombinedLogAppender">
        <!-- MongoDB配置 -->
        <server>127.0.0.1</server>
        <port>27017</port>
        <databaseName>nifi_logs</databaseName>
        <userName>nifi</userName>
        <password>xxxx</password>
        <collectionName>app_logs</collectionName>
        <hostIp>${HOSTIP}</hostIp>
    </appender>

    <logger name="org.apache.nifi" level="INFO" additivity="false">
        <appender-ref ref="CombinedLogAppender" />
    </logger>

    <root level="WARN">
        <appender-ref ref="CombinedLogAppender" />
    </root>
</configuration>

注意事项

  • OpenTelemetry默认会和团队部署的OTel Collector自动对接,无需额外配置ES地址,和其他应用的上报规则保持一致
  • 可根据团队ES索引规范调整自定义属性的key名,适配现有日志查询规则
  • 测试阶段可临时添加调试日志验证OTel上报链路,正常运行后可关闭调试日志避免性能损耗

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 23:30:02