Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -1103,11 +1103,17 @@ && compareWriterProgress(requestProgress, storedWriterProgress) == 0) {
"resolved first uncovered replayable WAL record");
}
return ReplayLocateDecision.atEnd(
consensusReqReader.getCurrentSearchIndex(),
nextSearchIndexAfterCurrent(),
effectiveRecoveryRegionProgress,
"all locally replayable WAL records are already covered");
}

private long nextSearchIndexAfterCurrent() {
// The reader reports the last local WAL index, while the subscription cursor is the next
// index to read. Starting at the last index would replay it and report one spurious WAL gap.
return consensusReqReader.getCurrentSearchIndex() + 1L;
}

protected ReplayLocateDecision locateReplayStartForRegionProgress(
final RegionProgress regionProgress, final boolean seekAfter) {
if (!(consensusReqReader instanceof WALNode)) {
Expand Down Expand Up @@ -3068,19 +3074,19 @@ public void cleanUp() {
// ======================== Seek ========================

/**
* Seeks to the earliest available WAL position. The actual position depends on WAL retention: if
* old files have been reclaimed, the earliest available position may be later than 0.
* Seeks to the earliest available WAL position. Local consensus search indexes start at 1; 0
* denotes an empty WAL. If old files have been reclaimed, replay may start at a later index.
*/
public void seekToBeginning() {
seekToResolvedPosition(0L, new RegionProgress(Collections.emptyMap()), "beginning");
seekToResolvedPosition(
FIRST_CONSENSUS_SEARCH_INDEX, new RegionProgress(Collections.emptyMap()), "beginning");
}

/**
* Seeks to the current WAL write position. After this, only newly written data will be consumed.
*/
public void seekToEnd() {
seekToResolvedPosition(
consensusReqReader.getCurrentSearchIndex(), computeTailRegionProgress(), "end");
seekToResolvedPosition(nextSearchIndexAfterCurrent(), computeTailRegionProgress(), "end");
}

public void seekToRegionProgress(final RegionProgress regionProgress) {
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,241 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/

package org.apache.iotdb.db.subscription.broker.consensus;

import org.apache.iotdb.commons.conf.CommonConfig;
import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.consensus.DataRegionId;
import org.apache.iotdb.consensus.common.request.IndexedConsensusRequest;
import org.apache.iotdb.consensus.iot.IoTConsensusServerImpl;
import org.apache.iotdb.consensus.iot.SubscriptionWalRetentionPolicy;
import org.apache.iotdb.consensus.iot.WriterSafeFrontierTracker;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode;
import org.apache.iotdb.db.queryengine.plan.statement.StatementTestUtils;
import org.apache.iotdb.db.storageengine.dataregion.wal.node.WALNode;
import org.apache.iotdb.db.subscription.event.SubscriptionEvent;
import org.apache.iotdb.db.subscription.resource.SubscriptionMemoryManager;
import org.apache.iotdb.db.subscription.task.execution.ConsensusSubscriptionPrefetchExecutor;
import org.apache.iotdb.db.subscription.task.execution.ConsensusSubscriptionPrefetchExecutorManager;
import org.apache.iotdb.rpc.subscription.config.TopicConstant;
import org.apache.iotdb.rpc.subscription.payload.poll.RegionProgress;
import org.apache.iotdb.rpc.subscription.payload.poll.SubscriptionPollResponseType;
import org.apache.iotdb.rpc.subscription.payload.poll.TabletsPayload;

import org.apache.tsfile.enums.ColumnCategory;
import org.apache.tsfile.enums.TSDataType;
import org.apache.tsfile.write.record.Tablet;
import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.TemporaryFolder;
import org.junit.runner.RunWith;
import org.powermock.api.mockito.PowerMockito;
import org.powermock.core.classloader.annotations.PowerMockIgnore;
import org.powermock.core.classloader.annotations.PrepareForTest;
import org.powermock.modules.junit4.PowerMockRunner;

import java.io.File;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.Iterator;
import java.util.List;
import java.util.concurrent.TimeUnit;

import static org.junit.Assert.assertEquals;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;

@RunWith(PowerMockRunner.class)
@PrepareForTest(ConsensusSubscriptionPrefetchExecutorManager.class)
@PowerMockIgnore({"com.sun.org.apache.xerces.*", "javax.xml.*", "org.xml.*", "javax.management.*"})
public class ConsensusPrefetchingQueueSeekTest {

@Rule public final TemporaryFolder temporaryFolder = new TemporaryFolder();

@Test(timeout = 10_000L)
public void testSeekToBeginningReplaysFirstLocalEntryWithoutWalGap() throws Exception {
assertBeginningReplay(Collections.singletonList(createRequest(1L, 1L, 7)), 2L, 0L);
}

@Test(timeout = 10_000L)
public void testSeekToBeginningPreservesFollowerEntryBeforeFirstLocalEntry() throws Exception {
assertBeginningReplay(
Arrays.asList(createRequest(-1L, 10L, 8), createRequest(1L, 1L, 7)), 2L, 0L);
}

@Test(timeout = 10_000L)
public void testSeekToBeginningCountsOnlyMissingValidLocalIndexes() throws Exception {
// Local indexes 1 and 2 are absent. Index 0 is the empty-WAL sentinel, not a missing entry.
assertBeginningReplay(
Arrays.asList(createRequest(3L, 3L, 7), createRequest(4L, 4L, 7)), 5L, 2L);
}

private void assertBeginningReplay(
final List<IndexedConsensusRequest> requests,
final long expectedNextIndex,
final long expectedSkippedEntries)
throws Exception {
final CommonConfig config = CommonDescriptor.getInstance().getConfig();
final int originalBatchMaxDelay = config.getSubscriptionConsensusBatchMaxDelayInMs();
final int originalPrefetchThreads =
config.getSubscriptionConsensusPrefetchExecutorMaxThreadNum();
final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();
final File systemDir = temporaryFolder.newFolder();
ConsensusSubscriptionPrefetchExecutor executor = null;
ConsensusPrefetchingQueue queue = null;
try {
config.setSubscriptionConsensusBatchMaxDelayInMs(0);
config.setSubscriptionConsensusPrefetchExecutorMaxThreadNum(1);
IoTDBDescriptor.getInstance().getConfig().setSystemDir(systemDir.getAbsolutePath());

// Supply an isolated runtime even in builds where subscriptions are disabled. Seek reset
// and replay still run through the real serial prefetch worker and the public queue API.
executor = new ConsensusSubscriptionPrefetchExecutor();
final ConsensusSubscriptionPrefetchExecutorManager manager =
mock(ConsensusSubscriptionPrefetchExecutorManager.class);
PowerMockito.mockStatic(ConsensusSubscriptionPrefetchExecutorManager.class);
PowerMockito.when(ConsensusSubscriptionPrefetchExecutorManager.getInstance())
.thenReturn(manager);
when(manager.getExecutor()).thenReturn(executor);

final WALNode walNode = mock(WALNode.class);
when(walNode.getCurrentSearchIndex()).thenReturn(expectedNextIndex - 1L);
when(walNode.getLogDirectory()).thenReturn(systemDir);
final IoTConsensusServerImpl server = mock(IoTConsensusServerImpl.class);
when(server.getConsensusReqReader()).thenReturn(walNode);
when(server.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker());
final ConsensusLogToTabletConverter converter = mock(ConsensusLogToTabletConverter.class);
when(converter.getDatabaseName()).thenReturn("db");
when(converter.isTableModel()).thenReturn(true);
when(converter.convert(any()))
.thenAnswer(
invocation ->
Collections.singletonList(
createTablet(((InsertRowNode) invocation.getArgument(0)).getTime())));

queue =
new ConsensusPrefetchingQueue(
"consumerGroup",
"topic",
TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
new DataRegionId(1),
server,
new SubscriptionWalRetentionPolicy(
"topic",
SubscriptionWalRetentionPolicy.UNBOUNDED,
SubscriptionWalRetentionPolicy.UNBOUNDED),
converter,
new ConsensusSubscriptionCommitManager(
(consumerGroupId, topicName, regionId) ->
ConsensusSubscriptionCommitManager.ConfigNodeProgressQueryResult.absent()),
new RegionProgress(Collections.emptyMap()),
expectedNextIndex,
1L,
true) {
@Override
protected ProgressWALIterator createSubscriptionWALIterator(
final long startSearchIndex) {
final Iterator<IndexedConsensusRequest> retainedEntries =
requests.stream()
.filter(
request ->
request.getSearchIndex() < 0
|| request.getSearchIndex() >= startSearchIndex)
.iterator();
final ProgressWALIterator iterator = mock(ProgressWALIterator.class);
when(iterator.hasNext()).thenAnswer(ignored -> retainedEntries.hasNext());
when(iterator.next()).thenAnswer(ignored -> retainedEntries.next());
return iterator;
}
};
queue.setSubscriptionMemoryManager(new SubscriptionMemoryManager(16L * 1024 * 1024));

queue.seekToBeginning();

final List<Long> actualTimestamps = new ArrayList<>();
final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(5L);
while (actualTimestamps.size() < requests.size() && System.nanoTime() < deadline) {
final SubscriptionEvent event = queue.poll("consumer");
if (event == null) {
TimeUnit.MILLISECONDS.sleep(10L);
continue;
}
assertEquals(
SubscriptionPollResponseType.TABLETS.getType(),
event.getCurrentResponse().getResponseType());
for (final Tablet tablet :
((TabletsPayload) event.getCurrentResponse().getPayload()).getTablets()) {
assertEquals(1, tablet.getRowSize());
actualTimestamps.add(tablet.getTimestamps()[0]);
}
}
final List<Long> expectedTimestamps = new ArrayList<>();
for (final IndexedConsensusRequest request : requests) {
expectedTimestamps.add(((InsertRowNode) request.getRequests().get(0)).getTime());
}
Collections.sort(expectedTimestamps);
Collections.sort(actualTimestamps);
assertEquals(expectedTimestamps, actualTimestamps);
assertEquals(expectedNextIndex, queue.getCurrentReadSearchIndex());
assertEquals(expectedSkippedEntries, queue.getWalGapSkippedEntries());
verify(converter, times(requests.size())).convert(any());
} finally {
if (queue != null) {
queue.close();
}
if (executor != null) {
executor.shutdown();
}
config.setSubscriptionConsensusBatchMaxDelayInMs(originalBatchMaxDelay);
config.setSubscriptionConsensusPrefetchExecutorMaxThreadNum(originalPrefetchThreads);
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
}
}

private static IndexedConsensusRequest createRequest(
final long searchIndex, final long localSeq, final int writerNodeId) {
return new IndexedConsensusRequest(
searchIndex,
localSeq,
Collections.singletonList(
StatementTestUtils.genInsertRowNode(Math.toIntExact(localSeq))))
.setPhysicalTime(1000L + localSeq)
.setNodeId(writerNodeId);
}

private static Tablet createTablet(final long timestamp) {
final Tablet tablet =
new Tablet(
"sensors",
Arrays.asList("device", "temperature"),
Arrays.asList(TSDataType.STRING, TSDataType.DOUBLE),
Arrays.asList(ColumnCategory.TAG, ColumnCategory.FIELD),
1);
tablet.addTimestamp(0, timestamp);
tablet.addValue(0, 0, "d1");
tablet.addValue(0, 1, 36.5);
tablet.setRowSize(1);
return tablet;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -612,12 +612,13 @@ public void testReplayStartPreservesUncoveredFollowerEntries() throws Exception
reader.currentSearchIndex = 5L;
final ConsensusPrefetchingQueue.ReplayLocateDecision tailDecision =
queue.scanReplayStartForRequests(
Collections.singletonList(createRequest(-1L, 11L, 101L, 8)).iterator(),
Arrays.asList(createRequest(5L, 10L, 100L, 8), createRequest(-1L, 11L, 101L, 8))
.iterator(),
regionProgress,
true);

assertEquals(ConsensusPrefetchingQueue.ReplayLocateStatus.AT_END, tailDecision.getStatus());
assertEquals(5L, tailDecision.getStartSearchIndex());
assertEquals(6L, tailDecision.getStartSearchIndex());
assertEquals(
committedProgress,
tailDecision.getRecoveryRegionProgress().getWriterPositions().get(formerLeader));
Expand Down Expand Up @@ -1129,6 +1130,7 @@ public void testAtEndReplayLookupPreservesRequestedWriterFrontier() throws Excep
queue.scanReplayStartForRequests(Collections.emptyIterator(), requestedProgress, true);

assertEquals(ConsensusPrefetchingQueue.ReplayLocateStatus.AT_END, decision.getStatus());
assertEquals(6L, decision.getStartSearchIndex());
assertEquals(requestedProgress, decision.getRecoveryRegionProgress());
} finally {
if (queue != null) {
Expand All @@ -1138,6 +1140,53 @@ public void testAtEndReplayLookupPreservesRequestedWriterFrontier() throws Excep
}
}

@Test
public void testAtEndInitializationHasNoTailGap() throws Exception {
final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();
final File systemDir = temporaryFolder.newFolder("atEndNoTailGap");
final File walDirectory = temporaryFolder.newFolder("atEndNoTailGapWal");
ConsensusPrefetchingQueue queue = null;
try {
final DataRegionId regionId = new DataRegionId(7);
final WALNode walNode = mock(WALNode.class);
when(walNode.getLogDirectory()).thenReturn(walDirectory);
when(walNode.getCurrentSearchIndex()).thenReturn(5L);
final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class);
when(serverImpl.getConsensusReqReader()).thenReturn(walNode);
when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker());
final RegionProgress committedProgress =
new RegionProgress(
Collections.singletonMap(
new WriterId(regionId.toString(), 7), new WriterProgress(1000L, 5L)));
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),
committedProgress,
6L,
1L,
true);

assertNull(queue.poll("consumer"));
assertEquals(6L, queue.getCurrentReadSearchIndex());
assertEquals(0L, queue.getRawWalGap());
} finally {
if (queue != null) {
queue.close();
}
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
}
}

@Test
public void testAtEndInitializationPreservesPendingQueueRegistrationBoundary() throws Exception {
final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();
Expand Down
Loading