diff --git a/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java b/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java index 8357af332..cea7c47a6 100644 --- a/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java +++ b/topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java @@ -17,6 +17,7 @@ */ public class MessageDecoder { private static final Logger logger = LoggerFactory.getLogger(MessageDecoder.class); + private final long maxBufferSize; private final AtomicLong totalAvailable; private final Executor decompressionExecutor; @@ -26,6 +27,11 @@ public class MessageDecoder { private volatile boolean isStopped = false; public MessageDecoder(long maxBufferSize, Executor decompressionExecutor, CodecRegistry codecRegistry) { + if (maxBufferSize <= 0) { + throw new IllegalArgumentException("maxBufferSize must be positive, but got " + maxBufferSize); + } + + this.maxBufferSize = maxBufferSize; this.totalAvailable = new AtomicLong(maxBufferSize); this.decompressionExecutor = decompressionExecutor; this.codecRegistry = codecRegistry; @@ -61,12 +67,26 @@ void free(long bufferSize) { private final class DecodeNext implements Runnable { @Override public void run() { - while (!isStopped && totalAvailable.get() > 0) { - ReadPartitionDecoder.EncodedMessage next = decodingQueue.poll(); + while (!isStopped) { + long available = totalAvailable.get(); + if (available <= 0) { + return; + } + + // Only this runnable polls the queue and it is serialized, so peek() cannot be invalidated here + ReadPartitionDecoder.EncodedMessage next = decodingQueue.peek(); if (next == null) { return; } + // A message larger than the whole budget must still make progress, but only when nothing else + // retains the buffer. Otherwise it waits until the already admitted messages are released. + if (next.getUncompressedSize() > available && available != maxBufferSize) { + return; + } + + decodingQueue.poll(); + long size = next.allocate(); totalAvailable.addAndGet(-size); try { diff --git a/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java b/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java index 3e526e21f..3da03cf98 100644 --- a/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java +++ b/topic/src/main/java/tech/ydb/topic/read/impl/ReadPartitionDecoder.java @@ -91,6 +91,14 @@ private void release() { decoder.free(allocatedTotal.getAndSet(0)); } + private void notifyReady() { + try { + readyHandler.run(); + } catch (Throwable th) { + logger.error("[{}] Exception was thrown by the ready handler", traceID, th); + } + } + public class EncodedMessage extends MessageImpl { private final int codecCode; private final long uncompressedSize; @@ -125,7 +133,11 @@ public void setError(Throwable th) { problem = new IOException("Decompression for " + getPartitionSession() + " error", th); releaseRange(OffsetsRange.of(getOffset())); isReady = true; - readyHandler.run(); + notifyReady(); + } + + long getUncompressedSize() { + return uncompressedSize; } public long allocate() { @@ -172,7 +184,7 @@ public void decode(CodecRegistry registry) { data = null; } isReady = true; - readyHandler.run(); + notifyReady(); } } } diff --git a/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java b/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java index 47ae01876..b997958bd 100644 --- a/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java +++ b/topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java @@ -80,6 +80,64 @@ private static BatchMeta meta(int codec) { @Rule public final HideLoggersRule hideLogger = new HideLoggersRule(); + @Test + public void nonPositiveBufferSizeTest() { + // A decoder with a non-positive budget can never admit a message and silently stalls the reader + Assert.assertThrows(IllegalArgumentException.class, + () -> new MessageDecoder(0, Runnable::run, REGISTRY)); + Assert.assertThrows(IllegalArgumentException.class, + () -> new MessageDecoder(-1, Runnable::run, REGISTRY)); + } + + @Test + @HideLoggers(MessageDecoder.class) + public void readyHandlerThrowsOnDecodeTest() { + MessageDecoder decoder = new MessageDecoder(1000, Runnable::run, REGISTRY); + + AtomicInteger ready = new AtomicInteger(0); + ReadPartitionDecoder partition = new ReadPartitionDecoder("t1", decoder, PS1, null, () -> { + ready.incrementAndGet(); + throw new RuntimeException("ready handler is broken"); + }); + + BatchMeta meta = meta(Codec.GZIP); + MessageImpl m1 = partition.decode(meta, OffsetsRange.of(1), gzipMsg(1, 40)); + MessageImpl m2 = partition.decode(meta, OffsetsRange.of(2), gzipMsg(2, 50)); + + decoder.decodeNext(); + + // A broken handler must not escape into the decompression thread and must not stop the following messages + Assert.assertEquals(2, ready.get()); + Assert.assertTrue(m1.isReady()); + Assert.assertTrue(m2.isReady()); + Assert.assertEquals(40, m1.getData().length); + Assert.assertEquals(50, m2.getData().length); + } + + @Test + @HideLoggers(MessageDecoder.class) + public void readyHandlerThrowsOnErrorTest() { + Executor rejecting = task -> { + throw new RejectedExecutionException("executor is saturated"); + }; + MessageDecoder decoder = new MessageDecoder(1000, rejecting, REGISTRY); + + AtomicInteger ready = new AtomicInteger(0); + ReadPartitionDecoder partition = new ReadPartitionDecoder("t1", decoder, PS1, null, () -> { + ready.incrementAndGet(); + throw new RuntimeException("ready handler is broken"); + }); + + BatchMeta meta = meta(Codec.GZIP); + MessageImpl m1 = partition.decode(meta, OffsetsRange.of(1), gzipMsg(1, 40)); + + decoder.decodeNext(); + + Assert.assertEquals(1, ready.get()); + Assert.assertTrue(m1.isReady()); + assertDecompressionException("Decompression for " + PS1 + " error", m1::getData); + } + @Test public void rawDecodeTest() { MessageDecoder decoder = new MessageDecoder(10000, Runnable::run, REGISTRY); @@ -130,45 +188,56 @@ public void flowControlByBudgetTest() { decoder.decodeNext(); - Assert.assertEquals(-50, decoder.getTotalAvailable()); + // 40 + 50 fit into the budget, 60 does not and has to wait + Assert.assertEquals(10, decoder.getTotalAvailable()); Assert.assertTrue(m1.isReady()); Assert.assertTrue(m2.isReady()); - Assert.assertTrue(m3.isReady()); + Assert.assertFalse(m3.isReady()); Assert.assertFalse(m4.isReady()); Assert.assertFalse(m5.isReady()); - Assert.assertEquals(3, ready.get()); + Assert.assertEquals(2, ready.get()); Assert.assertEquals(40, m1.getData().length); Assert.assertEquals(50, m2.getData().length); - Assert.assertEquals(60, m3.getData().length); Assert.assertEquals(1, m1.getData()[0]); Assert.assertEquals(2, m2.getData()[0]); - Assert.assertEquals(3, m3.getData()[0]); - p1.releaseRange(OffsetsRange.of(1)); // 40 is not enough to resume decoding + p1.releaseRange(OffsetsRange.of(1)); // 10 + 40 is not enough for the 60 bytes of m3 p1.releaseRange(OffsetsRange.of(4)); // that offset is not decoded yet - Assert.assertEquals(-10, decoder.getTotalAvailable()); + Assert.assertEquals(50, decoder.getTotalAvailable()); + Assert.assertFalse(m3.isReady()); Assert.assertFalse(m4.isReady()); Assert.assertFalse(m5.isReady()); - Assert.assertEquals(3, ready.get()); + Assert.assertEquals(2, ready.get()); - p1.releaseRange(OffsetsRange.of(2)); - Assert.assertEquals(-30, decoder.getTotalAvailable()); - Assert.assertTrue(m4.isReady()); + p1.releaseRange(OffsetsRange.of(2)); // the whole budget is free again, m3 is admitted + Assert.assertEquals(40, decoder.getTotalAvailable()); + Assert.assertTrue(m3.isReady()); + Assert.assertFalse(m4.isReady()); Assert.assertFalse(m5.isReady()); - Assert.assertEquals(4, ready.get()); + Assert.assertEquals(3, ready.get()); - Assert.assertEquals(70, m4.getData().length); - Assert.assertEquals(4, m4.getData()[0]); + Assert.assertEquals(60, m3.getData().length); + Assert.assertEquals(3, m3.getData()[0]); p1.releaseRange(OffsetsRange.of(0, 3)); // double release - Assert.assertEquals(-30, decoder.getTotalAvailable()); + Assert.assertEquals(40, decoder.getTotalAvailable()); + Assert.assertFalse(m4.isReady()); + Assert.assertEquals(3, ready.get()); + + p1.releaseRange(OffsetsRange.of(0, 5)); // releases m3, then m4 fits into the free budget + Assert.assertEquals(30, decoder.getTotalAvailable()); + Assert.assertTrue(m4.isReady()); Assert.assertFalse(m5.isReady()); Assert.assertEquals(4, ready.get()); - p1.releaseRange(OffsetsRange.of(0, 5)); + Assert.assertEquals(70, m4.getData().length); + Assert.assertEquals(4, m4.getData()[0]); + + // m5 is bigger than the whole budget, so it is admitted only when nothing else retains the buffer + p1.releaseRange(OffsetsRange.of(0, 6)); Assert.assertTrue(m5.isReady()); Assert.assertEquals(5, ready.get()); @@ -199,7 +268,8 @@ public void partitionFlowTest() { Assert.assertEquals(100, decoder.getTotalAvailable()); decoder.decodeNext(); - Assert.assertEquals(-50, decoder.getTotalAvailable()); + // 40 + 50 fit into the budget, the 60 bytes of m3 do not + Assert.assertEquals(10, decoder.getTotalAvailable()); MessageImpl m6 = p1.decode(meta, OffsetsRange.of(4), gzipMsg(4, 10)); MessageImpl m7 = p1.decode(meta, OffsetsRange.of(5), gzipMsg(5, 20)); @@ -207,11 +277,12 @@ public void partitionFlowTest() { MessageImpl m9 = p2.decode(meta, OffsetsRange.of(15), gzipMsg(15, 20)); decoder.decodeNext(); - Assert.assertEquals(-50, decoder.getTotalAvailable()); + // m3 still blocks the queue, the smaller messages behind it are not reordered + Assert.assertEquals(10, decoder.getTotalAvailable()); Assert.assertTrue(m1.isReady()); Assert.assertTrue(m2.isReady()); - Assert.assertTrue(m3.isReady()); + Assert.assertFalse(m3.isReady()); Assert.assertFalse(m4.isReady()); Assert.assertFalse(m5.isReady()); Assert.assertFalse(m6.isReady()); @@ -220,8 +291,9 @@ public void partitionFlowTest() { Assert.assertFalse(m9.isReady()); Assert.assertEquals(2, r1.get()); - Assert.assertEquals(1, r2.get()); + Assert.assertEquals(0, r2.get()); + // closing p1 returns its 90 bytes, which is enough to admit m3 and then m4 p1.close(); Assert.assertEquals(0, decoder.getTotalAvailable()); @@ -270,23 +342,27 @@ public void decodeStopTest() { decoder.decodeNext(); - Assert.assertEquals(-20, decoder.getTotalAvailable()); - Assert.assertEquals(2, ready.get()); + // only m1 fits into the 70 bytes budget, the 50 bytes of m2 do not + Assert.assertEquals(30, decoder.getTotalAvailable()); + Assert.assertEquals(1, ready.get()); Assert.assertTrue(m1.isReady()); - Assert.assertTrue(m2.isReady()); + Assert.assertFalse(m2.isReady()); Assert.assertFalse(m3.isReady()); decoder.stop(); - Assert.assertEquals(-20, decoder.getTotalAvailable()); + Assert.assertEquals(30, decoder.getTotalAvailable()); - Assert.assertEquals(2, ready.get()); + Assert.assertEquals(1, ready.get()); + Assert.assertFalse(m2.isReady()); Assert.assertFalse(m3.isReady()); + // a stopped decoder neither returns the budget nor resumes the pending messages partition.releaseRange(OffsetsRange.of(0, 10)); - Assert.assertEquals(2, ready.get()); + Assert.assertEquals(1, ready.get()); + Assert.assertFalse(m2.isReady()); Assert.assertFalse(m3.isReady()); - Assert.assertEquals(-20, decoder.getTotalAvailable()); + Assert.assertEquals(30, decoder.getTotalAvailable()); } @Test @@ -319,7 +395,8 @@ public void decodesOnProvidedExecutorTest() { Assert.assertFalse(m1.isReady()); Assert.assertFalse(m2.isReady()); - Assert.assertEquals(3, decodeTasks.size()); + // 400 + 500 fit into the 1000 bytes budget, the 600 bytes of m3 do not + Assert.assertEquals(2, decodeTasks.size()); Assert.assertEquals(0, p1ready.get()); Assert.assertEquals(0, p2ready.get()); @@ -327,7 +404,7 @@ public void decodesOnProvidedExecutorTest() { Assert.assertTrue(m1.isReady()); Assert.assertFalse(m2.isReady()); - Assert.assertEquals(2, decodeTasks.size()); + Assert.assertEquals(1, decodeTasks.size()); Assert.assertEquals(1, p1ready.get()); Assert.assertEquals(0, p2ready.get()); @@ -335,18 +412,30 @@ public void decodesOnProvidedExecutorTest() { Assert.assertTrue(m1.isReady()); Assert.assertTrue(m2.isReady()); - Assert.assertEquals(1, decodeTasks.size()); + Assert.assertEquals(0, decodeTasks.size()); Assert.assertEquals(1, p1ready.get()); Assert.assertEquals(1, p2ready.get()); + // decoding a message does not return its budget, only releasing it does + p1.releaseRange(OffsetsRange.of(0, 2)); + Assert.assertEquals(0, decodeTasks.size()); + + // now the whole budget is free again and m3 and m4 are admitted + p2.releaseRange(OffsetsRange.of(0, 2)); + Assert.assertEquals(2, decodeTasks.size()); + decoder.stop(); p1.close(); p2.close(); - Assert.assertEquals(1, decodeTasks.size()); + // tasks already submitted to the executor must not decode after the partitions are closed + decodeTasks.poll().run(); decodeTasks.poll().run(); Assert.assertEquals(0, decodeTasks.size()); + Assert.assertFalse(m3.isReady()); + Assert.assertFalse(m4.isReady()); + Assert.assertFalse(m5.isReady()); Assert.assertEquals(1, p1ready.get()); Assert.assertEquals(1, p2ready.get()); }