Caideyipi commented on code in PR #18298:
URL: https://github.com/apache/iotdb/pull/18298#discussion_r3662763829


##########
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:
   Added detailed, localized messages for invalid component counts, negative 
values, and a non-positive total, with message assertions in 
DataNodeMemoryConfigTest.



##########
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:
   Merged the two branches since both block real-time admission and return the 
same idle result.



##########
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:
   Renamed it to requestWithEmptyConversionResult. The request itself contains 
data; only the converter's second result is empty.



##########
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:
   Replaced the arbitrary 20 with backpressureVerificationRounds = writeCount. 
This provides one scheduling attempt per submitted write and verifies repeated 
prefetch cannot bypass backpressure.



-- 
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]

Reply via email to