Skip to content

Commit e49d9ad

Browse files
committed
Merge branch 'master' of https://github.com/apache/iotdb into revert-silent
2 parents 0901540 + ce14b31 commit e49d9ad

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
@@ -1492,6 +1492,10 @@ public PrefetchRoundResult drivePrefetchOnce() {
14921492
}
14931493
if (batchResult != MaterializationResult.SUCCESS) {
14941494
if (batchResult == MaterializationResult.WAL_GAP) {
1495+
if (!lingerBatch.isEmpty() && !flushBatch(lingerBatch, observedSeekGeneration)) {
1496+
resetRoundStateForSeek(seekGeneration.get());
1497+
return PrefetchRoundResult.rescheduleNow();
1498+
}
14951499
return PrefetchRoundResult.rescheduleAfter(WAL_GAP_RETRY_SLEEP_MS);
14961500
}
14971501
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
@@ -1523,6 +1523,94 @@ public void testWalReplayCountsOnlyUnavailableSearchIndexes() throws Exception {
15231523
}
15241524
}
15251525

1526+
@Test
1527+
public void testPendingWalGapFlushesLingerBatchBeforeRetry() throws Exception {
1528+
final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();
1529+
final int originalBatchMaxWalEntries =
1530+
CommonDescriptor.getInstance().getConfig().getSubscriptionConsensusBatchMaxWalEntries();
1531+
final int originalBatchMaxDelay =
1532+
CommonDescriptor.getInstance().getConfig().getSubscriptionConsensusBatchMaxDelayInMs();
1533+
final File systemDir = temporaryFolder.newFolder("pending-gap-flush-linger-batch");
1534+
ConsensusPrefetchingQueue queue = null;
1535+
try {
1536+
CommonDescriptor.getInstance().getConfig().setSubscriptionConsensusBatchMaxWalEntries(2);
1537+
CommonDescriptor.getInstance().getConfig().setSubscriptionConsensusBatchMaxDelayInMs(60_000);
1538+
1539+
final DataRegionId regionId = new DataRegionId(13);
1540+
final WALNode walNode = mock(WALNode.class);
1541+
when(walNode.getCurrentSearchIndex()).thenReturn(5L);
1542+
when(walNode.getLogDirectory()).thenReturn(systemDir);
1543+
final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class);
1544+
when(serverImpl.getConsensusReqReader()).thenReturn(walNode);
1545+
when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker());
1546+
1547+
final ConsensusLogToTabletConverter converter = mock(ConsensusLogToTabletConverter.class);
1548+
when(converter.convert(any())).thenReturn(Collections.singletonList(createTablet()));
1549+
when(converter.getDatabaseName()).thenReturn("db");
1550+
1551+
queue =
1552+
new ConsensusPrefetchingQueue(
1553+
"consumerGroup",
1554+
"topic",
1555+
TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
1556+
regionId,
1557+
serverImpl,
1558+
new SubscriptionWalRetentionPolicy(
1559+
"topic",
1560+
SubscriptionWalRetentionPolicy.UNBOUNDED,
1561+
SubscriptionWalRetentionPolicy.UNBOUNDED),
1562+
converter,
1563+
newCommitManager(systemDir),
1564+
new RegionProgress(Collections.emptyMap()),
1565+
1L,
1566+
1L,
1567+
true) {
1568+
@Override
1569+
protected ProgressWALIterator createSubscriptionWALIterator(
1570+
final long startSearchIndex) {
1571+
final Iterator<IndexedConsensusRequest> walEntries =
1572+
Arrays.asList(
1573+
createRequest(1L),
1574+
createRequest(2L),
1575+
createRequest(3L),
1576+
createRequest(4L),
1577+
createRequest(5L))
1578+
.iterator();
1579+
final ProgressWALIterator iterator = mock(ProgressWALIterator.class);
1580+
when(iterator.hasNext()).thenAnswer(ignored -> walEntries.hasNext());
1581+
when(iterator.next()).thenAnswer(ignored -> walEntries.next());
1582+
return iterator;
1583+
}
1584+
};
1585+
queue.setSubscriptionMemoryManager(new SubscriptionMemoryManager(16L * 1024 * 1024));
1586+
1587+
assertNull(queue.poll("consumer"));
1588+
assertTrue(pendingEntries(queue).offer(createRequest(5L)));
1589+
1590+
queue.drivePrefetchOnce();
1591+
1592+
assertEquals(3L, queue.getCurrentReadSearchIndex());
1593+
assertEquals(1, queue.getPrefetchedEventCount());
1594+
final SubscriptionEvent event = queue.poll("consumer");
1595+
assertNotNull(event);
1596+
assertEquals(
1597+
SubscriptionPollResponseType.TABLETS.getType(),
1598+
event.getCurrentResponse().getResponseType());
1599+
assertTrue(queue.ack("consumer", event.getCommitContext()));
1600+
} finally {
1601+
if (queue != null) {
1602+
queue.close();
1603+
}
1604+
CommonDescriptor.getInstance()
1605+
.getConfig()
1606+
.setSubscriptionConsensusBatchMaxWalEntries(originalBatchMaxWalEntries);
1607+
CommonDescriptor.getInstance()
1608+
.getConfig()
1609+
.setSubscriptionConsensusBatchMaxDelayInMs(originalBatchMaxDelay);
1610+
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
1611+
}
1612+
}
1613+
15261614
@Test
15271615
public void testActivationRetriesUntilConfigNodeProgressIsExplicitlyAvailable() throws Exception {
15281616
final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();

0 commit comments

Comments
 (0)