I am using apache Nifi , which uses logback framework for logging. I have written a custom log appender to push logs to MongoDB. Now I have a requirement to push logs in specific format (json) to ElasticSearch.
As few other apps are already using open telemetry to push data to ES, i have been asked to do so.
So now I am looking for right way to push data from my custom log appender (java class) to Open Telemetry.
Could some please point me to the usefull example which I can refer ?
Mongo appender looks like below, same message I want to push Open Telemetry also.
@Override
public void start() {
super.start();
logger.info("Initialising MongoDBLogAppender ~~~~~~~~~~~~~~~~");
try{
createConnection();
}catch (Exception e){
logger.error("Failed to obtain a database instance from the MongoClient at server [{}] and port [{}].", 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 {
logger.error("Mongo database name is required.");
}
return client;
}
@Override
protected void append(ILoggingEvent iLoggingEvent) {
if (database == null)
return;
String logMessage = iLoggingEvent.getMessage() == null ? "" : iLoggingEvent.getMessage();
logger.info("LogMessage (inside append method) : "+logMessage);//Change the level to debug after testing
String[] msgParts = parseLogMessage(logMessage);
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); // <--- Data requested and the insertion will now occur
}
@Override
public void onNext(final InsertOneResult result) {}
@Override
public void onError(final Throwable t) {
logger.error("Failed to insert Nifi log to mongodb : "+this.toString(),t);
}
@Override
public void onComplete() {}
});
}
catch (Exception e)
{
logger.error("Encountered exception while logging Nifi log to MongoDB : "+this.toString(), e);
}
}
/* Log message format expected : "~(<LoggedinTenantId>, <TargetTenantId>, <Userid>) ~ <RequestId> ~ <SessionId> ~ <DetailedMessage>" */
public static String[] parseLogMessage(String logMessage){
String[] msgParts = new String[6];
..................................
return msgParts;
}
private synchronized static MongoClient createMongoClient(String server, int port, String databaseName, String userName, String password)
{
...................................
return MongoClients.create(settings);
}
}
Thanks Mahendra