Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 22 additions & 2 deletions topic/src/main/java/tech/ydb/topic/read/impl/MessageDecoder.java
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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() {
Expand Down Expand Up @@ -172,7 +184,7 @@ public void decode(CodecRegistry registry) {
data = null;
}
isReady = true;
readyHandler.run();
notifyReady();
}
}
}
Expand Down
151 changes: 120 additions & 31 deletions topic/src/test/java/tech/ydb/topic/read/impl/MessageDecoderTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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());

Expand Down Expand Up @@ -199,19 +268,21 @@ 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));
MessageImpl m8 = p2.decode(meta, OffsetsRange.of(14), gzipMsg(14, 10));
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());
Expand All @@ -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());

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -319,34 +395,47 @@ 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());

decodeTasks.poll().run();

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());

decodeTasks.poll().run();

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());
}
Expand Down
Loading