jt2594838 commented on code in PR #18298:
URL: https://github.com/apache/iotdb/pull/18298#discussion_r3655570615
##########
iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java:
##########
@@ -749,6 +857,550 @@ public void
testActivationInstallsRuntimeStateBeforeRefreshingAuthoritativeProgr
}
}
+ @Test
+ public void testTabletMemoryReleasedAfterAck() throws Exception {
+ final String originalSystemDir =
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
+ final File systemDir = temporaryFolder.newFolder("system-memory-release");
+ ConsensusPrefetchingQueue queue = null;
+ try {
+ final DataRegionId regionId = new DataRegionId(1);
+ final FakeConsensusReqReader reader = new FakeConsensusReqReader();
+ final IoTConsensusServerImpl serverImpl =
mock(IoTConsensusServerImpl.class);
+ when(serverImpl.getConsensusReqReader()).thenReturn(reader);
+ when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new
WriterSafeFrontierTracker());
+
+ final ConsensusLogToTabletConverter converter =
mock(ConsensusLogToTabletConverter.class);
+ when(converter.convert(any()))
+ .thenReturn(Collections.singletonList(createTablet()),
Collections.emptyList());
+ when(converter.getDatabaseName()).thenReturn("db");
+
+ queue =
+ new ConsensusPrefetchingQueue(
+ "consumerGroup",
+ "topic",
+ TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
+ regionId,
+ serverImpl,
+ new SubscriptionWalRetentionPolicy(
+ "topic",
+ SubscriptionWalRetentionPolicy.UNBOUNDED,
+ SubscriptionWalRetentionPolicy.UNBOUNDED),
+ converter,
+ newCommitManager(systemDir),
+ new RegionProgress(Collections.emptyMap()),
+ 1L,
+ 1L,
+ true);
+
+ final IndexedConsensusRequest dataRequest =
+ new IndexedConsensusRequest(
+ 1L,
Collections.singletonList(StatementTestUtils.genInsertRowNode(1)))
+ .setPhysicalTime(1000L)
+ .setNodeId(7);
+ final IndexedConsensusRequest emptyRequest =
+ new IndexedConsensusRequest(
+ 2L,
Collections.singletonList(StatementTestUtils.genInsertRowNode(2)))
+ .setPhysicalTime(1001L)
+ .setNodeId(7);
+ reader.currentSearchIndex = 2L;
+ pendingEntries(queue).offer(dataRequest);
+ pendingEntries(queue).offer(emptyRequest);
Review Comment:
Why is the second one called emptyRequest?
##########
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/conf/DataNodeMemoryConfig.java:
##########
@@ -464,6 +474,42 @@ private void initCompactionMemoryManager(long
compactionMemorySize) {
"Compaction", compactionMemorySize - fixedMemoryCost);
}
+ static int[] resolveQueryMemoryProportions(
+ final String configuredProportions, final boolean subscriptionEnabled) {
+ final int[] resolvedProportions = DEFAULT_QUERY_MEMORY_PROPORTIONS.clone();
+ if (configuredProportions != null) {
+ final String[] proportions = configuredProportions.split(":");
+ if (proportions.length != LEGACY_QUERY_MEMORY_COMPONENT_COUNT
+ && proportions.length != QUERY_MEMORY_COMPONENT_COUNT) {
+ throw new IllegalArgumentException();
Review Comment:
Give some detailed message?
##########
iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java:
##########
@@ -749,6 +857,550 @@ public void
testActivationInstallsRuntimeStateBeforeRefreshingAuthoritativeProgr
}
}
+ @Test
+ public void testTabletMemoryReleasedAfterAck() throws Exception {
+ final String originalSystemDir =
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
+ final File systemDir = temporaryFolder.newFolder("system-memory-release");
+ ConsensusPrefetchingQueue queue = null;
+ try {
+ final DataRegionId regionId = new DataRegionId(1);
+ final FakeConsensusReqReader reader = new FakeConsensusReqReader();
+ final IoTConsensusServerImpl serverImpl =
mock(IoTConsensusServerImpl.class);
+ when(serverImpl.getConsensusReqReader()).thenReturn(reader);
+ when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new
WriterSafeFrontierTracker());
+
+ final ConsensusLogToTabletConverter converter =
mock(ConsensusLogToTabletConverter.class);
+ when(converter.convert(any()))
+ .thenReturn(Collections.singletonList(createTablet()),
Collections.emptyList());
+ when(converter.getDatabaseName()).thenReturn("db");
+
+ queue =
+ new ConsensusPrefetchingQueue(
+ "consumerGroup",
+ "topic",
+ TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
+ regionId,
+ serverImpl,
+ new SubscriptionWalRetentionPolicy(
+ "topic",
+ SubscriptionWalRetentionPolicy.UNBOUNDED,
+ SubscriptionWalRetentionPolicy.UNBOUNDED),
+ converter,
+ newCommitManager(systemDir),
+ new RegionProgress(Collections.emptyMap()),
+ 1L,
+ 1L,
+ true);
+
+ final IndexedConsensusRequest dataRequest =
+ new IndexedConsensusRequest(
+ 1L,
Collections.singletonList(StatementTestUtils.genInsertRowNode(1)))
+ .setPhysicalTime(1000L)
+ .setNodeId(7);
+ final IndexedConsensusRequest emptyRequest =
+ new IndexedConsensusRequest(
+ 2L,
Collections.singletonList(StatementTestUtils.genInsertRowNode(2)))
+ .setPhysicalTime(1001L)
+ .setNodeId(7);
+ reader.currentSearchIndex = 2L;
+ pendingEntries(queue).offer(dataRequest);
+ pendingEntries(queue).offer(emptyRequest);
+
+ assertNull(queue.poll("consumer"));
+ queue.drivePrefetchOnce();
+
+ final long retainedBytes = queue.getRetainedTabletBytes();
+ assertTrue(retainedBytes > 0L);
+ final SubscriptionEvent event = queue.poll("consumer");
+ assertNotNull(event);
+ assertEquals(retainedBytes, queue.getRetainedTabletBytes());
+
+ assertTrue(queue.ack("consumer", event.getCommitContext()));
+ assertEquals(0L, queue.getRetainedTabletBytes());
+ } finally {
+ if (queue != null) {
+ queue.close();
+ }
+
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
+ }
+ }
+
+ @Test
+ public void testCleanupReconcilesUnindexedTabletReservation() throws
Exception {
+ final String originalSystemDir =
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
+ final File systemDir =
temporaryFolder.newFolder("system-orphan-memory-release");
+ ConsensusPrefetchingQueue queue = null;
+ try {
+ final DataRegionId regionId = new DataRegionId(1);
+ final FakeConsensusReqReader reader = new FakeConsensusReqReader();
+ final IoTConsensusServerImpl serverImpl =
mock(IoTConsensusServerImpl.class);
+ when(serverImpl.getConsensusReqReader()).thenReturn(reader);
+ when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new
WriterSafeFrontierTracker());
+
+ queue =
+ new ConsensusPrefetchingQueue(
+ "consumerGroup",
+ "topic",
+ TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
+ regionId,
+ serverImpl,
+ new SubscriptionWalRetentionPolicy(
+ "topic",
+ SubscriptionWalRetentionPolicy.UNBOUNDED,
+ SubscriptionWalRetentionPolicy.UNBOUNDED),
+ mock(ConsensusLogToTabletConverter.class),
+ newCommitManager(systemDir),
+ new RegionProgress(Collections.emptyMap()),
+ 1L,
+ 1L,
+ true);
+ final SubscriptionMemoryManager memoryManager = new
SubscriptionMemoryManager(1024L);
+ queue.setSubscriptionMemoryManager(memoryManager);
+
+ final Method reserveTabletMemory =
+
ConsensusPrefetchingQueue.class.getDeclaredMethod("tryReserveTabletMemory",
long.class);
+ reserveTabletMemory.setAccessible(true);
+ assertTrue((Boolean) reserveTabletMemory.invoke(queue, 256L));
+ assertEquals(256L, queue.getRetainedTabletBytes());
+ assertEquals(256L, memoryManager.getUsedMemorySizeInBytes());
+
+ queue.cleanUp();
+
+ assertEquals(0L, queue.getRetainedTabletBytes());
+ assertEquals(0L, memoryManager.getUsedMemorySizeInBytes());
+ } finally {
+ if (queue != null) {
+ queue.close();
+ }
+
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
+ }
+ }
+
+ @Test
+ public void testMemoryBackpressureStopsFurtherMaterializationUntilAck()
throws Exception {
+ final String originalSystemDir =
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
+ final File systemDir =
temporaryFolder.newFolder("system-memory-backpressure");
+ ConsensusPrefetchingQueue queue = null;
+ try {
+ final DataRegionId regionId = new DataRegionId(1);
+ final FakeConsensusReqReader reader = new FakeConsensusReqReader();
+ final IoTConsensusServerImpl serverImpl =
mock(IoTConsensusServerImpl.class);
+ when(serverImpl.getConsensusReqReader()).thenReturn(reader);
+ when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new
WriterSafeFrontierTracker());
+
+ final AtomicInteger conversionCount = new AtomicInteger();
+ final ConsensusLogToTabletConverter converter =
mock(ConsensusLogToTabletConverter.class);
+ when(converter.convert(any()))
+ .thenAnswer(
+ ignored -> {
+ conversionCount.incrementAndGet();
+ return Collections.singletonList(createTablet());
+ });
+ when(converter.getDatabaseName()).thenReturn("db");
+
+ queue =
+ new ConsensusPrefetchingQueue(
+ "consumerGroup",
+ "topic",
+ TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
+ regionId,
+ serverImpl,
+ new SubscriptionWalRetentionPolicy(
+ "topic",
+ SubscriptionWalRetentionPolicy.UNBOUNDED,
+ SubscriptionWalRetentionPolicy.UNBOUNDED),
+ converter,
+ newCommitManager(systemDir),
+ new RegionProgress(Collections.emptyMap()),
+ 1L,
+ 1L,
+ true);
+ final long oneTabletBytes = createTablet().ramBytesUsed();
+ final SubscriptionMemoryManager memoryManager =
+ new SubscriptionMemoryManager(oneTabletBytes + Math.max(1L,
oneTabletBytes / 2L));
+ queue.setSubscriptionMemoryManager(memoryManager);
+
+ final IndexedConsensusRequest firstRequest = createRequest(1L);
+ final IndexedConsensusRequest secondRequest = createRequest(2L);
+ final IndexedConsensusRequest thirdRequest = createRequest(3L);
+ assertTrue(pendingEntries(queue).offer(firstRequest));
+ assertTrue(pendingEntries(queue).offer(secondRequest));
+ assertTrue(pendingEntries(queue).offer(thirdRequest));
+ reader.currentSearchIndex = 3L;
+
+ assertNull(queue.poll("consumer"));
+ queue.drivePrefetchOnce();
+
+ assertEquals(2, conversionCount.get());
+ assertEquals(2L, queue.getCurrentReadSearchIndex());
+ assertEquals(oneTabletBytes, queue.getRetainedTabletBytes());
+ assertTrue(memoryManager.getFreeMemorySizeInBytes() > 0L);
+ assertEquals(1, queue.getPrefetchedEventCount());
+ assertTrue(pendingEntries(queue).isEmpty());
+ assertEquals("true",
queue.coreReportMessage().get("realtimeAdmissionBlocked"));
+ assertFalse(pendingEntries(queue).offer(createRequest(4L)));
+
+ queue.drivePrefetchOnce();
+ assertEquals(2, conversionCount.get());
+ assertEquals(oneTabletBytes, queue.getRetainedTabletBytes());
+
+ final SubscriptionEvent event = queue.poll("consumer");
+ assertNotNull(event);
+ assertTrue(queue.ack("consumer", event.getCommitContext()));
+ assertEquals(0L, queue.getRetainedTabletBytes());
+
+ queue.drivePrefetchOnce();
+ assertEquals("false",
queue.coreReportMessage().get("realtimeAdmissionBlocked"));
+ assertTrue(pendingEntries(queue).offer(secondRequest));
+ assertTrue(pendingEntries(queue).offer(thirdRequest));
+ queue.drivePrefetchOnce();
+ assertEquals(4, conversionCount.get());
+ assertEquals(3L, queue.getCurrentReadSearchIndex());
+ assertEquals(oneTabletBytes, queue.getRetainedTabletBytes());
+ assertEquals(1, queue.getPrefetchedEventCount());
+
+ final SubscriptionEvent secondEvent = queue.poll("consumer");
+ assertNotNull(secondEvent);
+ assertTrue(queue.ack("consumer", secondEvent.getCommitContext()));
+ queue.drivePrefetchOnce();
+ assertTrue(pendingEntries(queue).offer(thirdRequest));
+ queue.drivePrefetchOnce();
+
+ assertEquals(5, conversionCount.get());
+ assertEquals(4L, queue.getCurrentReadSearchIndex());
+ assertEquals("0", queue.coreReportMessage().get("pendingEntriesSize"));
+ assertEquals("0",
queue.coreReportMessage().get("bufferedRealtimeEntryCount"));
+ assertEquals(oneTabletBytes, queue.getRetainedTabletBytes());
+
+ queue.close();
+ queue = null;
+ assertEquals(0L, memoryManager.getUsedMemorySizeInBytes());
+ } finally {
+ if (queue != null) {
+ queue.close();
+ }
+
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
+ }
+ }
+
+ @Test
+ public void testWideTablePausedConsumerKeepsMaterializedMemoryBounded()
throws Exception {
+ final String originalSystemDir =
IoTDBDescriptor.getInstance().getConfig().getSystemDir();
+ final File systemDir =
temporaryFolder.newFolder("system-wide-table-memory-bound");
+ ConsensusPrefetchingQueue queue = null;
+ try {
+ final DataRegionId regionId = new DataRegionId(1);
+ final FakeConsensusReqReader reader = new FakeConsensusReqReader();
+ final IoTConsensusServerImpl serverImpl =
mock(IoTConsensusServerImpl.class);
+ when(serverImpl.getConsensusReqReader()).thenReturn(reader);
+ when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new
WriterSafeFrontierTracker());
+
+ final int columnCount = 128;
+ final int rowCount = 64;
+ final long oneTabletBytes = createWideTablet(columnCount,
rowCount).ramBytesUsed();
+ final AtomicInteger conversionCount = new AtomicInteger();
+ final ConsensusLogToTabletConverter converter =
mock(ConsensusLogToTabletConverter.class);
+ when(converter.convert(any()))
+ .thenAnswer(
+ ignored -> {
+ conversionCount.incrementAndGet();
+ return Collections.singletonList(createWideTablet(columnCount,
rowCount));
+ });
+ when(converter.getDatabaseName()).thenReturn("db");
+
+ queue =
+ new ConsensusPrefetchingQueue(
+ "consumerGroup",
+ "topic",
+ TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
+ regionId,
+ serverImpl,
+ new SubscriptionWalRetentionPolicy(
+ "topic",
+ SubscriptionWalRetentionPolicy.UNBOUNDED,
+ SubscriptionWalRetentionPolicy.UNBOUNDED),
+ converter,
+ newCommitManager(systemDir),
+ new RegionProgress(Collections.emptyMap()),
+ 1L,
+ 1L,
+ true);
+ final long memoryLimit = oneTabletBytes * 2L + oneTabletBytes / 2L;
+ final SubscriptionMemoryManager memoryManager = new
SubscriptionMemoryManager(memoryLimit);
+ queue.setSubscriptionMemoryManager(memoryManager);
+
+ final int writeCount = 256;
+ for (long searchIndex = 1L; searchIndex <= writeCount; searchIndex++) {
+ assertTrue(pendingEntries(queue).offer(createRequest(searchIndex)));
+ }
+ reader.currentSearchIndex = writeCount;
+
+ assertNull(queue.poll("pausedConsumer"));
+ queue.drivePrefetchOnce();
+
+ assertEquals(3, conversionCount.get());
+ assertEquals(3L, queue.getCurrentReadSearchIndex());
+ assertEquals(oneTabletBytes * 2L, queue.getRetainedTabletBytes());
+ assertTrue(queue.getRetainedTabletBytes() <=
queue.getSubscriptionMemoryLimitInBytes());
+ assertEquals(1, queue.getPrefetchedEventCount());
+ assertEquals("0", queue.coreReportMessage().get("pendingEntriesSize"));
+ assertEquals("0",
queue.coreReportMessage().get("bufferedRealtimeEntryCount"));
+ assertEquals("true",
queue.coreReportMessage().get("realtimeAdmissionBlocked"));
+
+ for (int round = 0; round < 20; round++) {
+ queue.drivePrefetchOnce();
+ }
Review Comment:
Why prefetch 20 rounds?
##########
iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java:
##########
@@ -1258,6 +1357,31 @@ public PrefetchRoundResult drivePrefetchOnce() {
final int maxTablets =
config.getSubscriptionConsensusBatchMaxTabletCount();
final long maxBatchBytes =
config.getSubscriptionConsensusBatchMaxSizeInBytes();
+ // Always consume already materialized entries before converting more
WAL requests into
+ // Tablets. Otherwise a full delivery batch can leave an unbounded
hidden writer backlog.
+ if (!drainBufferedRealtimeWriters(
+ lingerBatch, observedSeekGeneration, maxTablets, maxBatchBytes)) {
+ resetRoundStateForSeek(seekGeneration.get());
+ return PrefetchRoundResult.rescheduleNow();
+ }
+ if (prefetchingQueue.size() >= MAX_PREFETCHING_QUEUE_SIZE) {
+ blockRealtimeAdmission();
+ return computeIdleRoundResult();
+ }
+ if (!realtimeEntriesByWriter.isEmpty()) {
+ blockRealtimeAdmission();
+ return computeIdleRoundResult();
+ }
Review Comment:
Merge branches?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]