Skip to content

Commit 7662d3f

Browse files
committed
Log improved
1 parent 37cd0eb commit 7662d3f

4 files changed

Lines changed: 55 additions & 21 deletions

File tree

kafka/kafka/src/main/java/io/koraframework/kafka/common/consumer/containers/KafkaAssignConsumerContainer.java

Lines changed: 21 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -301,11 +301,15 @@ private Consumer<K, V> initializeConsumer() {
301301
}
302302
return new ConsumerWrapper<>(realConsumer, driverMetrics, keyDeserializer, valueDeserializer);
303303
} catch (Exception e) {
304-
logger.error("KafkaListener '{}' failed to start in assign mode, due to: {}", listenerName, e.getMessage(), e);
304+
logger.atError()
305+
.addKeyValue("listenerName", this.listenerName)
306+
.log("KafkaListener failed to start in assign mode, due to: {}", e.getMessage(), e);
305307
try {
306308
Thread.sleep(250);
307309
} catch (InterruptedException ie) {
308-
logger.error("KafkaListener '{}' error interrupting thread", listenerName, ie);
310+
logger.atError()
311+
.addKeyValue("listenerName", this.listenerName)
312+
.log("KafkaListener error interrupting thread", ie);
309313
}
310314
return null;
311315
}
@@ -316,7 +320,9 @@ public void init() {
316320
var threads = this.threads;
317321
if (threads > 0) {
318322
if (this.topic != null) {
319-
logger.debug("KafkaListener '{}' starting in assign mode...", listenerName);
323+
logger.atDebug()
324+
.addKeyValue("listenerName", this.listenerName)
325+
.log("KafkaListener starting in assign mode...", listenerName);
320326
final long started = TimeUtils.started();
321327

322328
executorService = Executors.newFixedThreadPool(threads, new NamedThreadFactory(listenerName));
@@ -355,7 +361,9 @@ public void init() {
355361
@Override
356362
public void release() {
357363
if (isActive.compareAndSet(true, false)) {
358-
logger.debug("KafkaListener '{}' stopping...", listenerName);
364+
logger.atDebug()
365+
.addKeyValue("listenerName", this.listenerName)
366+
.log("KafkaListener stopping...");
359367
var started = System.nanoTime();
360368

361369
for (var consumer : consumers) {
@@ -364,11 +372,15 @@ public void release() {
364372
consumers.clear();
365373
if (executorService != null) {
366374
if (!shutdownExecutorService(executorService, config.shutdownWait())) {
367-
logger.warn("KafkaListener '{}' failed completing graceful shutdown in {}", listenerName, config.shutdownWait());
375+
logger.atWarn()
376+
.addKeyValue("listenerName", this.listenerName)
377+
.log("KafkaListener failed completing graceful shutdown in {}", config.shutdownWait());
368378
}
369379
}
370380

371-
logger.info("KafkaListener '{}' stopped in {}", listenerName, TimeUtils.tookForLogging(started));
381+
logger.atInfo()
382+
.addKeyValue("listenerName", this.listenerName)
383+
.log("KafkaListener stopped in {}", TimeUtils.tookForLogging(started));
372384
}
373385
}
374386

@@ -377,7 +389,9 @@ private boolean shutdownExecutorService(ExecutorService executorService, Duratio
377389
if (!terminated) {
378390
executorService.shutdown();
379391
try {
380-
logger.debug("KafkaListener '{}' awaiting graceful shutdown...", listenerName);
392+
logger.atDebug()
393+
.addKeyValue("listenerName", this.listenerName)
394+
.log("KafkaListener awaiting graceful shutdown...");
381395
terminated = executorService.awaitTermination(shutdownAwait.toMillis(), TimeUnit.MILLISECONDS);
382396
if (!terminated) {
383397
executorService.shutdownNow();

kafka/kafka/src/main/java/io/koraframework/kafka/common/consumer/containers/KafkaSubscribeConsumerContainer.java

Lines changed: 21 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -146,7 +146,9 @@ public void launchPollLoop(Consumer<K, V> consumer, String listenerLogName, long
146146
@Override
147147
public void init() {
148148
if (config.threads() > 0 && this.isActive.compareAndSet(false, true)) {
149-
logger.debug("KafkaListener '{}' starting in subscribe mode...", listenerName);
149+
logger.atDebug()
150+
.addKeyValue("listenerName", this.listenerName)
151+
.log("KafkaListener starting in subscribe mode...");
150152
final long started = TimeUtils.started();
151153

152154
executorService = Executors.newFixedThreadPool(config.threads(), new NamedThreadFactory(listenerName));
@@ -184,7 +186,9 @@ public void init() {
184186
@Override
185187
public void release() {
186188
if (isActive.compareAndSet(true, false)) {
187-
logger.debug("KafkaListener '{}' stopping...", listenerName);
189+
logger.atDebug()
190+
.addKeyValue("listenerName", this.listenerName)
191+
.log("KafkaListener stopping...");
188192
final long started = TimeUtils.started();
189193

190194
for (var consumer : consumers) {
@@ -193,12 +197,15 @@ public void release() {
193197
consumers.clear();
194198
if (executorService != null) {
195199
if (!shutdownExecutorService(executorService, config.shutdownWait())) {
196-
logger.warn("KafkaListener '{}' failed completing graceful shutdown in {}",
197-
listenerName, config.shutdownWait());
200+
logger.atWarn()
201+
.addKeyValue("listenerName", this.listenerName)
202+
.log("KafkaListener failed completing graceful shutdown in {}", config.shutdownWait());
198203
}
199204
}
200205

201-
logger.info("KafkaListener '{}' stopped in {}", listenerName, TimeUtils.tookForLogging(started));
206+
logger.atInfo()
207+
.addKeyValue("listenerName", this.listenerName)
208+
.log("KafkaListener stopped in {}", TimeUtils.tookForLogging(started));
202209
}
203210
}
204211

@@ -207,7 +214,9 @@ private boolean shutdownExecutorService(ExecutorService executorService, Duratio
207214
if (!terminated) {
208215
executorService.shutdown();
209216
try {
210-
logger.debug("KafkaListener '{}' awaiting graceful shutdown...", listenerName);
217+
logger.atDebug()
218+
.addKeyValue("listenerName", this.listenerName)
219+
.log("KafkaListener awaiting graceful shutdown...");
211220
terminated = executorService.awaitTermination(shutdownAwait.toMillis(), TimeUnit.MILLISECONDS);
212221
if (!terminated) {
213222
executorService.shutdownNow();
@@ -227,12 +236,15 @@ private Consumer<K, V> initializeConsumer() {
227236
try {
228237
return this.buildConsumer();
229238
} catch (Exception e) {
230-
logger.error("KafkaListener '{}' failed to start in subscribe mode, due to: {}",
231-
listenerName, e.getMessage(), e);
239+
logger.atError()
240+
.addKeyValue("listenerName", this.listenerName)
241+
.log("KafkaListener failed to start in subscribe mode, due to: {}", e.getMessage(), e);
232242
try {
233243
Thread.sleep(250);
234244
} catch (InterruptedException ie) {
235-
logger.error("KafkaListener '{}' error interrupting thread", listenerName, ie);
245+
logger.atError()
246+
.addKeyValue("listenerName", this.listenerName)
247+
.log("KafkaListener error interrupting thread", ie);
236248
}
237249
return null;
238250
}

kafka/kafka/src/main/java/io/koraframework/kafka/common/producer/AbstractPublisher.java

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,9 @@ public void init() throws Exception {
5353
return;
5454
}
5555
try {
56-
logger.debug("KafkaPublisher '{}' starting...", publisherName);
56+
logger.atDebug()
57+
.addKeyValue("publisherName", this.publisherName)
58+
.log("KafkaPublisher starting...");
5759
final long started = TimeUtils.started();
5860

5961
var producer = new KafkaProducer<>(driverProperties, new ByteArraySerializer(), new ByteArraySerializer());
@@ -63,7 +65,9 @@ public void init() throws Exception {
6365
this.micrometerMetrics.bindTo(this.telemetry.meterRegistry());
6466
}
6567

66-
logger.debug("KafkaPublisher '{}' started in {}", publisherName, TimeUtils.tookForLogging(started));
68+
logger.atInfo()
69+
.addKeyValue("publisherName", this.publisherName)
70+
.log("KafkaPublisher started in {}", TimeUtils.tookForLogging(started));
6771
} catch (Exception e) {
6872
throw new RuntimeException("KafkaPublisher '" + publisherName + "' failed to start, due to: " + e.getMessage(), e);
6973
}
@@ -75,7 +79,9 @@ public void release() throws Exception {
7579
return;
7680
}
7781

78-
logger.debug("KafkaPublisher '{}' stopping...", publisherName);
82+
logger.atDebug()
83+
.addKeyValue("publisherName", this.publisherName)
84+
.log("KafkaPublisher stopping...");
7985
final long started = TimeUtils.started();
8086

8187
var delegate = this.delegate;
@@ -84,6 +90,8 @@ public void release() throws Exception {
8490
this.micrometerMetrics = null;
8591
try (delegate; micrometerMetrics) {}
8692

87-
logger.info("KafkaPublisher '{}' stopped in {}", publisherName, TimeUtils.tookForLogging(started));
93+
logger.atInfo()
94+
.addKeyValue("publisherName", this.publisherName)
95+
.log("KafkaPublisher stopped in {}", TimeUtils.tookForLogging(started));
8896
}
8997
}

kafka/kafka/src/main/java/io/koraframework/kafka/common/producer/TransactionalPublisherImpl.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,8 @@
11
package io.koraframework.kafka.common.producer;
22

3+
import io.koraframework.application.graph.Lifecycle;
34
import org.apache.kafka.common.KafkaException;
45
import org.apache.kafka.common.errors.TimeoutException;
5-
import io.koraframework.application.graph.Lifecycle;
66

77
import java.util.Objects;
88
import java.util.concurrent.BlockingDeque;

0 commit comments

Comments
 (0)