From 17b39e75cdc5fae27ad2e5dbda2622d5fda2a0cb Mon Sep 17 00:00:00 2001 From: Sebastiaan de Schaetzen Date: Fri, 10 Apr 2026 15:43:30 +0200 Subject: [PATCH 1/2] =?UTF-8?q?refactor:=20simplify=20eventSize=20method?= =?UTF-8?q?=20signatures=20=E2=9C=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/main/java/ca/pjer/logback/Worker.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/main/java/ca/pjer/logback/Worker.java b/src/main/java/ca/pjer/logback/Worker.java index 2927b1f..6ee1568 100644 --- a/src/main/java/ca/pjer/logback/Worker.java +++ b/src/main/java/ca/pjer/logback/Worker.java @@ -48,11 +48,11 @@ InputLogEvent asInputLogEvent(ILoggingEvent event) { private static final int EVENT_SIZE_PADDING = 26; private static final Charset EVENT_SIZE_CHARSET = Charset.forName("UTF-8"); - static final int eventSize(InputLogEvent event) { + static int eventSize(InputLogEvent event) { return eventSize(event.message()); } - static final int eventSize(String message) { + static int eventSize(String message) { return message.getBytes(EVENT_SIZE_CHARSET).length + EVENT_SIZE_PADDING; } From 6ea16d03104661a5dfb90c1487eb6000e750fb63 Mon Sep 17 00:00:00 2001 From: Sebastiaan de Schaetzen Date: Fri, 10 Apr 2026 16:00:10 +0200 Subject: [PATCH 2/2] =?UTF-8?q?refactor:=20replace=20InputLogEvent=20with?= =?UTF-8?q?=20MemoizedEvent=20=F0=9F=93=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit This caches calls to eventSize(), which in our case were consuming a significant amount of CPU usage. When you're logging many large log lines, one can quickly exceed the MAX_BATCH_SIZE limit. At this point, it will calculate the size of each log event twice (once as part of drainBatchFromQueue->batchSize, and once as drainBatchFromQueue->eventSize). The eventSize method internally calls `String#getBytes`, which calls `String#encodeUTF8`. For large log lines, this call is quite expensive. Calling it twice is just unnecessary. --- .../java/ca/pjer/logback/AsyncWorker.java | 31 ++++++++-------- .../java/ca/pjer/logback/MemoizedEvent.java | 35 +++++++++++++++++++ 2 files changed, 52 insertions(+), 14 deletions(-) create mode 100644 src/main/java/ca/pjer/logback/MemoizedEvent.java diff --git a/src/main/java/ca/pjer/logback/AsyncWorker.java b/src/main/java/ca/pjer/logback/AsyncWorker.java index 3b4b410..a58f246 100644 --- a/src/main/java/ca/pjer/logback/AsyncWorker.java +++ b/src/main/java/ca/pjer/logback/AsyncWorker.java @@ -2,9 +2,7 @@ import static ca.pjer.logback.AwsLogsAppender.MAX_BATCH_LOG_EVENTS; -import java.util.ArrayDeque; -import java.util.Collection; -import java.util.Deque; +import java.util.*; import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; import java.util.concurrent.TimeUnit; @@ -22,7 +20,7 @@ class AsyncWorker extends Worker implements Runnable { private final int maxBatchLogEvents; private final int discardThreshold; private final AtomicBoolean running; - private final BlockingQueue queue; + private final BlockingQueue queue; private final AtomicLong lostCount; private Thread thread; @@ -32,7 +30,7 @@ class AsyncWorker extends Worker implements Runnable { maxBatchLogEvents = awsLogsAppender.getMaxBatchLogEvents(); discardThreshold = (int) Math.ceil(maxBatchLogEvents * 1.5); running = new AtomicBoolean(false); - queue = new ArrayBlockingQueue(maxBatchLogEvents * 2); + queue = new ArrayBlockingQueue<>(maxBatchLogEvents * 2); lostCount = new AtomicLong(0); } @@ -87,7 +85,7 @@ public void append(ILoggingEvent event) { long now = System.currentTimeMillis(); while (now < until) { try { - if (!queue.offer(logEvent, until - now, TimeUnit.MILLISECONDS)) { + if (!queue.offer(new MemoizedEvent(logEvent), until - now, TimeUnit.MILLISECONDS)) { lostCount.incrementAndGet(); AwsLogsMetricsHolder.get().incrementLostCount(); } @@ -104,7 +102,7 @@ public void append(ILoggingEvent event) { } } else { // we are not allowed to block, offer without blocking - if (!queue.offer(logEvent)) { + if (!queue.offer(new MemoizedEvent(logEvent))) { lostCount.incrementAndGet(); AwsLogsMetricsHolder.get().incrementLostCount(); } @@ -160,12 +158,12 @@ private void flush(boolean all) { private static final int MAX_BATCH_SIZE = 1048576; private Collection drainBatchFromQueue() { - Deque batch = new ArrayDeque(maxBatchLogEvents); + Deque batch = new ArrayDeque<>(maxBatchLogEvents); queue.drainTo(batch, MAX_BATCH_LOG_EVENTS); int batchSize = batchSize(batch); while (batchSize > MAX_BATCH_SIZE) { - InputLogEvent removed = batch.removeLast(); - batchSize -= eventSize(removed); + MemoizedEvent removed = batch.removeLast(); + batchSize -= removed.eventSize(); if (!queue.offer(removed)) { AwsLogsMetricsHolder.get().incrementBatchRequeueFailed(); if (getAwsLogsAppender().getVerbose()) { @@ -175,13 +173,18 @@ private Collection drainBatchFromQueue() { } AwsLogsMetricsHolder.get().incrementBatch(batchSize); - return batch; + + List array = new ArrayList<>(batchSize); + for (MemoizedEvent event : batch) { + array.add(event.event()); + } + return array; } - private static int batchSize(Collection batch) { + private static int batchSize(Collection batch) { int size = 0; - for (InputLogEvent event : batch) { - size += eventSize(event); + for (MemoizedEvent event : batch) { + size += event.eventSize(); } return size; } diff --git a/src/main/java/ca/pjer/logback/MemoizedEvent.java b/src/main/java/ca/pjer/logback/MemoizedEvent.java new file mode 100644 index 0000000..0424b5d --- /dev/null +++ b/src/main/java/ca/pjer/logback/MemoizedEvent.java @@ -0,0 +1,35 @@ +package ca.pjer.logback; + +import jdk.internal.RequiresIdentity; +import software.amazon.awssdk.services.cloudwatchlogs.model.InputLogEvent; + +import java.nio.charset.Charset; +import java.nio.charset.StandardCharsets; + +public class MemoizedEvent { + // See http://docs.aws.amazon.com/AmazonCloudWatchLogs/latest/APIReference/API_PutLogEvents.html + private static final int EVENT_SIZE_PADDING = 26; + private static final Charset EVENT_SIZE_CHARSET = StandardCharsets.UTF_8; + + private final InputLogEvent event; + private byte[] messageBytes; + + public MemoizedEvent(InputLogEvent event) { + this.event = event; + } + + public InputLogEvent event() { + return event; + } + + public byte[] getBytes() { + if (messageBytes == null) { + messageBytes = event.message().getBytes(EVENT_SIZE_CHARSET); + } + return messageBytes; + } + + public int eventSize() { + return getBytes().length + EVENT_SIZE_PADDING; + } +}