Skip to content

Commit ce14b31

Browse files
authored
[Subscription] Flush lingering batch before WAL gap retry (#18785)
1 parent 0c9fb83 commit ce14b31

2 files changed

Lines changed: 92 additions & 0 deletions

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1507,6 +1507,10 @@ public PrefetchRoundResult drivePrefetchOnce() {
15071507
}
15081508
if (batchResult != MaterializationResult.SUCCESS) {
15091509
if (batchResult == MaterializationResult.WAL_GAP) {
1510+
if (!lingerBatch.isEmpty() && !flushBatch(lingerBatch, observedSeekGeneration)) {
1511+
resetRoundStateForSeek(seekGeneration.get());
1512+
return PrefetchRoundResult.rescheduleNow();
1513+
}
15101514
return PrefetchRoundResult.rescheduleAfter(WAL_GAP_RETRY_SLEEP_MS);
15111515
}
15121516
if (batchResult == MaterializationResult.MEMORY_BLOCKED) {

‎iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java‎

Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1609,6 +1609,94 @@ public void testWalReplayCountsOnlyUnavailableSearchIndexes() throws Exception {
16091609
}
16101610
}
16111611

1612+
@Test
1613+
public void testPendingWalGapFlushesLingerBatchBeforeRetry() throws Exception {
1614+
final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();
1615+
final int originalBatchMaxWalEntries =
1616+
CommonDescriptor.getInstance().getConfig().getSubscriptionConsensusBatchMaxWalEntries();
1617+
final int originalBatchMaxDelay =
1618+
CommonDescriptor.getInstance().getConfig().getSubscriptionConsensusBatchMaxDelayInMs();
1619+
final File systemDir = temporaryFolder.newFolder("pending-gap-flush-linger-batch");
1620+
ConsensusPrefetchingQueue queue = null;
1621+
try {
1622+
CommonDescriptor.getInstance().getConfig().setSubscriptionConsensusBatchMaxWalEntries(2);
1623+
CommonDescriptor.getInstance().getConfig().setSubscriptionConsensusBatchMaxDelayInMs(60_000);
1624+
1625+
final DataRegionId regionId = new DataRegionId(13);
1626+
final WALNode walNode = mock(WALNode.class);
1627+
when(walNode.getCurrentSearchIndex()).thenReturn(5L);
1628+
when(walNode.getLogDirectory()).thenReturn(systemDir);
1629+
final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class);
1630+
when(serverImpl.getConsensusReqReader()).thenReturn(walNode);
1631+
when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker());
1632+
1633+
final ConsensusLogToTabletConverter converter = mock(ConsensusLogToTabletConverter.class);
1634+
when(converter.convert(any())).thenReturn(Collections.singletonList(createTablet()));
1635+
when(converter.getDatabaseName()).thenReturn("db");
1636+
1637+
queue =
1638+
new ConsensusPrefetchingQueue(
1639+
"consumerGroup",
1640+
"topic",
1641+
TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
1642+
regionId,
1643+
serverImpl,
1644+
new SubscriptionWalRetentionPolicy(
1645+
"topic",
1646+
SubscriptionWalRetentionPolicy.UNBOUNDED,
1647+
SubscriptionWalRetentionPolicy.UNBOUNDED),
1648+
converter,
1649+
newCommitManager(systemDir),
1650+
new RegionProgress(Collections.emptyMap()),
1651+
1L,
1652+
1L,
1653+
true) {
1654+
@Override
1655+
protected ProgressWALIterator createSubscriptionWALIterator(
1656+
final long startSearchIndex) {
1657+
final Iterator<IndexedConsensusRequest> walEntries =
1658+
Arrays.asList(
1659+
createRequest(1L),
1660+
createRequest(2L),
1661+
createRequest(3L),
1662+
createRequest(4L),
1663+
createRequest(5L))
1664+
.iterator();
1665+
final ProgressWALIterator iterator = mock(ProgressWALIterator.class);
1666+
when(iterator.hasNext()).thenAnswer(ignored -> walEntries.hasNext());
1667+
when(iterator.next()).thenAnswer(ignored -> walEntries.next());
1668+
return iterator;
1669+
}
1670+
};
1671+
queue.setSubscriptionMemoryManager(new SubscriptionMemoryManager(16L * 1024 * 1024));
1672+
1673+
assertNull(queue.poll("consumer"));
1674+
assertTrue(pendingEntries(queue).offer(createRequest(5L)));
1675+
1676+
queue.drivePrefetchOnce();
1677+
1678+
assertEquals(3L, queue.getCurrentReadSearchIndex());
1679+
assertEquals(1, queue.getPrefetchedEventCount());
1680+
final SubscriptionEvent event = queue.poll("consumer");
1681+
assertNotNull(event);
1682+
assertEquals(
1683+
SubscriptionPollResponseType.TABLETS.getType(),
1684+
event.getCurrentResponse().getResponseType());
1685+
assertTrue(queue.ack("consumer", event.getCommitContext()));
1686+
} finally {
1687+
if (queue != null) {
1688+
queue.close();
1689+
}
1690+
CommonDescriptor.getInstance()
1691+
.getConfig()
1692+
.setSubscriptionConsensusBatchMaxWalEntries(originalBatchMaxWalEntries);
1693+
CommonDescriptor.getInstance()
1694+
.getConfig()
1695+
.setSubscriptionConsensusBatchMaxDelayInMs(originalBatchMaxDelay);
1696+
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
1697+
}
1698+
}
1699+
16121700
@Test
16131701
public void testWalReplayFailsCriticallyWhenGapRemainsUnavailable() throws Exception {
16141702
final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();

0 commit comments

Comments
 (0)