Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,7 @@ private void close(Status status) {
if (shutdownFuture.isDone()) {
return;
}
config.getMetrics().close();

controlEventsExecutor.execute(() -> {
try {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,10 @@ public void init(String sessionId) {
requestFunc.accept(maxBufferSize);
}

long getCreditBalanceBytes() {
return maxBufferSize + totalReleased.get() - totalAllocated.get();
}

// Has no reentrant thread safety
public void allocate(long bufferSize, List<YdbTopic.StreamReadMessage.ReadResponse.PartitionData> dataList) {
logger.debug("[{}] Received ReadResponse of {} bytes, {} allocated and {} released before",
Expand Down
8 changes: 8 additions & 0 deletions topic/src/main/java/tech/ydb/topic/read/impl/ReadSession.java
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,14 @@ BufferManager getBufferManager() {
return bufferManager;
}

int getPartitionSessionCount() {
return partitions.size();
}

boolean isClosed() {
return isClosed;
}

ReadConfig getConfig() {
return config;
}
Expand Down
3 changes: 3 additions & 0 deletions topic/src/main/java/tech/ydb/topic/read/impl/ReaderImpl.java
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,7 @@ protected void onRetry(ReadSession stream, Status status) {
if (stream != null) {
stream.closeAll().forEach(ps -> handler.handleClosePartitionSession(ps));
}
config.getMetrics().unregisterStream(stream);
}

@Override
Expand All @@ -119,6 +120,7 @@ protected void onClose(ReadSession stream, Status status) {
if (stream != null) {
stream.closeAll().forEach(ps -> handler.handleClosePartitionSession(ps));
}
config.getMetrics().unregisterStream(stream);
handler.handleReaderClosed(status);
}

Expand All @@ -131,6 +133,7 @@ protected void onNext(ReadSession stream, FromServer message) {
currentSessionId = message.getInitResponse().getSessionId();
handler.handleSessionStarted(message.getInitResponse().getSessionId());
stream.onInit(message.getInitResponse());
config.getMetrics().registerStream(stream);
} else if (message.hasStartPartitionSessionRequest()) {
StartPartitionSessionEvent event = stream.onStartPartition(message.getStartPartitionSessionRequest());
if (event != null) {
Expand Down
50 changes: 46 additions & 4 deletions topic/src/main/java/tech/ydb/topic/read/impl/ReaderMetrics.java
Original file line number Diff line number Diff line change
Expand Up @@ -5,32 +5,74 @@
import tech.ydb.core.metrics.Attr;
import tech.ydb.core.metrics.LongCounter;
import tech.ydb.core.metrics.Meter;
import tech.ydb.core.metrics.MetricRegistration;

/**
* Topic reader counters.
* Topic reader metrics.
*/
final class ReaderMetrics {
final class ReaderMetrics implements AutoCloseable {
private static final String MESSAGE_UNIT = "{message}";

private final LongCounter deliveredMessages;
private final LongCounter receivedMessages;
private final LongCounter receivedBytes;
private final Attr[] commonAttributes;
private final boolean enabled;
private final Meter meter;
private MetricRegistration partitionsGauge = MetricRegistration.NOOP;
private MetricRegistration creditGauge = MetricRegistration.NOOP;
private ReadSession registeredStream;
private boolean closed;

ReaderMetrics(Meter meter, String consumer, String readerName) {
this.meter = meter;
this.enabled = meter != Meter.NOOP;
this.deliveredMessages = meter.createCounter(
"ydb.topic.reader.delivered.messages",
MESSAGE_UNIT,
"The number of messages delivered by the SDK to application code.");
this.receivedMessages = meter.createCounter("ydb.topic.reader.received.messages", MESSAGE_UNIT,
"Messages accepted by the SDK for active partition sessions.");
"The number of messages accepted by the SDK for an active partition session.");
this.receivedBytes = meter.createCounter("ydb.topic.reader.received.bytes", "By",
"Bytes in received ReadResponse messages.");
"The protocol bytes_size received in read responses.");
this.commonAttributes = createCommonAttributes(consumer, readerName);
}

synchronized void registerStream(ReadSession stream) {
if (closed || stream.isClosed()) {
return;
}
closeGauges();
registeredStream = stream;
partitionsGauge = meter.registerLongGauge(
"ydb.topic.reader.partition_session.count", "{session}",
"The number of partition sessions currently in the reader session processing lifecycle.",
m -> m.record(stream.getPartitionSessionCount(), commonAttributes));
creditGauge = meter.registerLongGauge("ydb.topic.reader.credit_balance_bytes", "By",
"The protocol credit granted to the server and not yet consumed by read responses.",
m -> m.record(stream.getBufferManager().getCreditBalanceBytes(), commonAttributes));
}

synchronized void unregisterStream(ReadSession stream) {
if (registeredStream == stream) {
closeGauges();
}
}

@Override
public synchronized void close() {
closed = true;
closeGauges();
}

private void closeGauges() {
partitionsGauge.close();
creditGauge.close();
partitionsGauge = MetricRegistration.NOOP;
creditGauge = MetricRegistration.NOOP;
registeredStream = null;
}

void reportDelivered(long messages, String topic) {
report(deliveredMessages, messages, topic);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,7 @@ public void shutdown() {
}

private void close(Status status) {
config.getMetrics().close();
initFuture.completeExceptionally(new RuntimeException("Reader was closed with " + status));
shutdownFuture.complete(status);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,19 @@
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;

import org.junit.Assert;
import org.junit.Test;
import org.mockito.Mockito;

import tech.ydb.core.metrics.Attr;
import tech.ydb.core.metrics.LongCounter;
import tech.ydb.core.metrics.LongMeasurement;
import tech.ydb.core.metrics.Meter;
import tech.ydb.core.metrics.MetricRegistration;
import tech.ydb.topic.TopicClient;
import tech.ydb.topic.TopicRpc;
import tech.ydb.topic.description.Codec;
Expand All @@ -25,6 +29,8 @@ public class ReaderMetricsTest {
private static final String DELIVERED = "ydb.topic.reader.delivered.messages";
private static final String RECEIVED_MESSAGES = "ydb.topic.reader.received.messages";
private static final String RECEIVED_BYTES = "ydb.topic.reader.received.bytes";
private static final String PARTITIONS = "ydb.topic.reader.partition_session.count";
private static final String CREDIT = "ydb.topic.reader.credit_balance_bytes";

@Test
public void readerCountersIncrementOnReceive() throws InterruptedException {
Expand Down Expand Up @@ -66,9 +72,83 @@ public void readerCountersIncrementOnReceive() throws InterruptedException {
}
}

@Test
public void gaugesObservePartitionSessionsAndProtocolCredit() throws InterruptedException {
RecordingMeter meter = new RecordingMeter();
ReadStreamMock stream = new ReadStreamMock();
TopicRpc rpc = Mockito.mock(TopicRpc.class);
Mockito.when(rpc.getScheduler()).thenReturn(Mockito.mock(ScheduledExecutorService.class));
Mockito.when(rpc.readSession(Mockito.anyString())).thenReturn(stream).thenReturn(new ReadStreamMock());
TopicClient client = TopicClientImpl.newClient(rpc).build();
SyncReader reader = client.createSyncReader(ReaderSettings.newBuilder()
.addTopic("/topic").setConsumerName("consumer")
.setMaxMemoryUsageBytes(100).withMeter(meter, "reader").build());
try {
Assert.assertTrue(meter.gauges.isEmpty());
reader.init();
Assert.assertTrue(meter.gauges.isEmpty());
stream.responseInit("session");
stream.responseStartPartition("/topic", 42, 0);
Assert.assertEquals(1, meter.collect(PARTITIONS));
Assert.assertEquals(100, meter.collect(CREDIT));
reader.init();
Assert.assertEquals(1, meter.collect(PARTITIONS));
Assert.assertEquals(100, meter.collect(CREDIT));

stream.responseData(20).partition(1, 0).batch(Codec.RAW, new byte[]{1}).and().send();
Assert.assertEquals(80, meter.collect(CREDIT));
Assert.assertNotNull(reader.receive(1, TimeUnit.SECONDS));
Assert.assertEquals(100, meter.collect(CREDIT));
} finally {
reader.shutdown();
client.close();
}
Assert.assertTrue(meter.gauges.isEmpty());
}

@Test
public void staleStreamCleanupDoesNotRemoveCurrentGauges() {
RecordingMeter meter = new RecordingMeter();
ReaderMetrics metrics = new ReaderMetrics(meter, "consumer", "reader");
ReadSession first = Mockito.mock(ReadSession.class);
ReadSession second = Mockito.mock(ReadSession.class);
Mockito.when(first.getPartitionSessionCount()).thenReturn(1);
Mockito.when(second.getPartitionSessionCount()).thenReturn(2);
try {
metrics.registerStream(first);
metrics.registerStream(second);
metrics.unregisterStream(first);
Assert.assertEquals(2, meter.collect(PARTITIONS));

Mockito.when(first.isClosed()).thenReturn(true);
metrics.registerStream(first);
Assert.assertEquals(2, meter.collect(PARTITIONS));
metrics.unregisterStream(second);
Assert.assertTrue(meter.gauges.isEmpty());
} finally {
metrics.close();
}
metrics.registerStream(second);
Assert.assertTrue(meter.gauges.isEmpty());
}

private static class RecordingMeter implements Meter {
private final Map<String, AtomicLong> counters = new ConcurrentHashMap<>();
private final Map<String, Attr[]> attributes = new ConcurrentHashMap<>();
private final Map<String, Consumer<LongMeasurement>> gauges = new ConcurrentHashMap<>();

@Override
public MetricRegistration registerLongGauge(
String name, String unit, String description, Consumer<LongMeasurement> callback) {
gauges.put(name, callback);
return () -> gauges.remove(name);
}

long collect(String name) {
Long[] value = new Long[1];
gauges.get(name).accept((observed, attrs) -> value[0] = observed);
return value[0];
}

@Override
public LongCounter createCounter(String name, String unit, String description) {
Expand Down
Loading