如何通过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
相关产品推荐
相关产品推荐

