Logback custom appender to push logs to ElasticSearch via Open Telemetry

Viewed 312

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

0 Answers
Related