Skip to content
Open
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 @@ -20,6 +20,7 @@
package org.apache.iotdb.db.pipe.event.common.tsfile.container.scan;

import org.apache.iotdb.commons.exception.IllegalPathException;
import org.apache.iotdb.commons.exception.pipe.PipeRuntimeOutOfMemoryCriticalException;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
import org.apache.iotdb.commons.pipe.datastructure.pattern.PipePattern;
Expand Down Expand Up @@ -352,6 +353,10 @@ private Tablet getNextTablet() {
}
PipeTabletUtils.compactBitMaps(tablet);
return tablet;
} catch (final PipeRuntimeOutOfMemoryCriticalException e) {
// Keep the parser state so the caller can yield its parser slot and retry from the same
// unconsumed data after memory is available again.
throw e;
} catch (final Exception e) {
close();
throw new PipeException("Failed to get next tablet insertion event.", e);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@ public class PipeMemoryManager {
private final Map<PipeIdentity, ArrayDeque<PipeRegionIdentity>>
waitingTsFileParserRegionOrderByPipe = new HashMap<>();
private final ArrayDeque<PipeIdentity> waitingTsFileParserPipeOrder = new ArrayDeque<>();
private PipeIdentity lastAdmittedWaitingTsFileParserPipe;

// Only non-zero memory blocks will be added to this set.
private final Set<PipeMemoryBlock> allocatedBlocks = new HashSet<>();
Expand Down Expand Up @@ -177,7 +178,8 @@ public synchronized boolean tryReserveTsFileParserMemory(
final PipeIdentity pipeIdentity = new PipeIdentity(pipeName, creationTime);
final PipeRegionIdentity pipeRegionIdentity =
new PipeRegionIdentity(pipeIdentity, dataRegionId);
enqueueTsFileParserReservationRequest(pipeRegionIdentity, reservationKey);
final boolean wasRequestAlreadyWaiting =
enqueueTsFileParserReservationRequest(pipeRegionIdentity, reservationKey);

final int globalLimit = Math.max(1, PIPE_CONFIG.getPipeTsFileParserInFlightMaxNum());
final int perPipeRegionLimit =
Expand Down Expand Up @@ -212,6 +214,9 @@ public synchronized boolean tryReserveTsFileParserMemory(
}

removeTsFileParserReservationRequest(pipeRegionIdentity, reservationKey, true);
if (wasRequestAlreadyWaiting) {
lastAdmittedWaitingTsFileParserPipe = pipeIdentity;
}
reservedTsFileParserCount++;
reservedTsFileParserCountByPipe.merge(pipeIdentity, 1, Integer::sum);
reservedTsFileParserCountByPipeRegion.put(pipeRegionIdentity, reservedCountOfPipeRegion + 1);
Expand Down Expand Up @@ -262,10 +267,11 @@ public synchronized void releaseTsFileParserMemory(
reservedTsFileParserCountByPipe.put(pipeIdentity, reservedCountOfPipe - 1);
}
reservedTsFileParserCount--;
clearTsFileParserAdmissionCursorIfIdle();
notifyNextTsFileParserMemoryReservationInternal();
}

private void enqueueTsFileParserReservationRequest(
private boolean enqueueTsFileParserReservationRequest(
final PipeRegionIdentity pipeRegionIdentity,
final TsFileParserMemoryReservation reservationKey) {
final LinkedHashSet<TsFileParserMemoryReservation> requestsOfPipeRegion =
Expand All @@ -282,7 +288,7 @@ private void enqueueTsFileParserReservationRequest(
regionOrder.addLast(key);
return new LinkedHashSet<>();
});
requestsOfPipeRegion.add(reservationKey);
return !requestsOfPipeRegion.add(reservationKey);
}

public synchronized void notifyNextTsFileParserMemoryReservation() {
Expand Down Expand Up @@ -322,7 +328,14 @@ private void notifyNextTsFileParserMemoryReservationInternal() {

private PipeRegionIdentity getNextEligibleTsFileParserPipeRegion(
final int perPipeRegionLimit, final boolean requirePipeWithoutReservedParser) {
PipeRegionIdentity firstEligiblePipeRegion = null;
boolean hasVisitedLastAdmittedPipe = lastAdmittedWaitingTsFileParserPipe == null;
for (final PipeIdentity pipeIdentity : waitingTsFileParserPipeOrder) {
final boolean isLastAdmittedPipe = pipeIdentity.equals(lastAdmittedWaitingTsFileParserPipe);
if (isLastAdmittedPipe) {
hasVisitedLastAdmittedPipe = true;
}

// Under soft memory pressure, reserve the hard-threshold headroom for a pipe that has no
// parser yet. Otherwise a busy pipe at the queue head can block every pipe behind it.
if (requirePipeWithoutReservedParser
Expand All @@ -335,14 +348,32 @@ private PipeRegionIdentity getNextEligibleTsFileParserPipeRegion(
if (regionOrder == null) {
continue;
}
PipeRegionIdentity eligiblePipeRegion = null;
for (final PipeRegionIdentity pipeRegionIdentity : regionOrder) {
if (reservedTsFileParserCountByPipeRegion.getOrDefault(pipeRegionIdentity, 0)
< perPipeRegionLimit) {
return pipeRegionIdentity;
eligiblePipeRegion = pipeRegionIdentity;
break;
}
}
if (eligiblePipeRegion == null) {
continue;
}

if (firstEligiblePipeRegion == null) {
firstEligiblePipeRegion = eligiblePipeRegion;
}
if (hasVisitedLastAdmittedPipe && !isLastAdmittedPipe) {
return eligiblePipeRegion;
}
}
return firstEligiblePipeRegion;
}

private void clearTsFileParserAdmissionCursorIfIdle() {
if (reservedTsFileParserCount == 0 && waitingTsFileParserPipeOrder.isEmpty()) {
lastAdmittedWaitingTsFileParserPipe = null;
}
return null;
}

private void removeTsFileParserReservationRequest(
Expand All @@ -365,6 +396,9 @@ private void removeTsFileParserReservationRequest(
if (regionOrder.isEmpty()) {
waitingTsFileParserRegionOrderByPipe.remove(pipeIdentity);
waitingTsFileParserPipeOrder.remove(pipeIdentity);
if (!rotateAfterAdmission) {
clearTsFileParserAdmissionCursorIfIdle();
}
return;
}
}
Expand All @@ -376,6 +410,8 @@ private void removeTsFileParserReservationRequest(
if (rotateAfterAdmission) {
waitingTsFileParserPipeOrder.remove(pipeIdentity);
waitingTsFileParserPipeOrder.addLast(pipeIdentity);
} else {
clearTsFileParserAdmissionCursorIfIdle();
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -184,6 +184,47 @@ public void testScanContainerReleasesTabletMemoryAfterRawTabletGenerated() throw
}
}

@Test
public void testScanContainerKeepsIteratorOnOutOfMemory() throws Exception {
nonalignedTsFile =
TsFileGeneratorUtils.generateNonAlignedTsFile(
"nonaligned-retry-tablet-memory.tsfile", 1, 1, 10, 0, 100, 10, 10);

try (final TsFileInsertionScanDataContainer container =
new TsFileInsertionScanDataContainer(
nonalignedTsFile,
new PrefixPipePattern("root"),
Long.MIN_VALUE,
Long.MAX_VALUE,
null,
null,
false)) {
final AtomicInteger memoryUsageReadCount = new AtomicInteger(0);
replaceAllocatedTabletMemory(
container,
new PipeMemoryBlock(0) {
@Override
public long getMemoryUsageInBytes() {
if (memoryUsageReadCount.incrementAndGet() == 2) {
throw new PipeRuntimeOutOfMemoryCriticalException("expected oom");
}
return super.getMemoryUsageInBytes();
}
});

final Iterator<TabletInsertionEvent> iterator =
container.toTabletInsertionEvents().iterator();
final PipeRuntimeOutOfMemoryCriticalException exception =
Assert.assertThrows(PipeRuntimeOutOfMemoryCriticalException.class, iterator::next);
Assert.assertEquals("expected oom", exception.getMessage());

Assert.assertTrue(iterator.hasNext());
final TabletInsertionEvent event = iterator.next();
Assert.assertTrue(event instanceof PipeRawTabletInsertionEvent);
((PipeRawTabletInsertionEvent) event).clearReferenceCount(getClass().getName());
}
}

@Test
public void testConsumeTabletInsertionEventsWithRetryPreservesProgressOnOutOfMemory()
throws Exception {
Expand Down Expand Up @@ -1309,6 +1350,16 @@ private PipeMemoryBlock getAllocatedTabletMemory(final TsFileInsertionDataContai
return (PipeMemoryBlock) field.get(container);
}

private void replaceAllocatedTabletMemory(
final TsFileInsertionDataContainer container, final PipeMemoryBlock replacement)
throws NoSuchFieldException, IllegalAccessException {
final Field field =
TsFileInsertionDataContainer.class.getDeclaredField("allocatedMemoryBlockForTablet");
field.setAccessible(true);
((PipeMemoryBlock) field.get(container)).close();
field.set(container, replacement);
}

@SuppressWarnings("unchecked")
private AtomicReference<TsFileInsertionDataContainer> getDataContainer(
final PipeTsFileInsertionEvent event) throws NoSuchFieldException, IllegalAccessException {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -192,6 +192,43 @@ public void testPipeFairnessIsNotWeightedByRegionCount() {
Assert.assertTrue(tryAcquire(pipeARegion2));
}

@Test
public void testPipeFairnessSurvivesTransientSingleRegionQueueGap() {
commonConfig.setPipeTsFileParserInFlightMaxNum(1);
commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(1);

final Reservation blocker = new Reservation("blocker", 0);
final Reservation multiRegion1 = new Reservation("multi", 1, "1");
final Reservation multiRegion2 = new Reservation("multi", 1, "2");
final Reservation multiRegion3 = new Reservation("multi", 1, "3");
final Reservation singleFirst = new Reservation("single", 2, "1");
final Reservation singleSecond = new Reservation("single", 2, "1");

Assert.assertTrue(tryAcquire(blocker));
Assert.assertFalse(tryAcquire(multiRegion1));
Assert.assertFalse(tryAcquire(multiRegion2));
Assert.assertFalse(tryAcquire(multiRegion3));
Assert.assertFalse(tryAcquire(singleFirst));

release(blocker);
Assert.assertTrue(tryAcquire(multiRegion1));
release(multiRegion1);

Assert.assertTrue(tryAcquire(singleFirst));
release(singleFirst);

// The single-region pipe temporarily has no admission request while it advances to its next
// TsFile, so the multi-region pipe can use the otherwise idle parser slot.
Assert.assertTrue(tryAcquire(multiRegion2));
Assert.assertFalse(tryAcquire(singleSecond));
release(multiRegion2);

// Once the single-region pipe is waiting again, the remembered pipe-level cursor must prevent
// another region of the multi-region pipe from taking a second consecutive turn.
Assert.assertFalse(tryAcquire(multiRegion3));
Assert.assertTrue(tryAcquire(singleSecond));
}

@Test
public void testSoftMemoryHeadroomIsReservedForPipeWithoutParser() {
commonConfig.setPipeTsFileParserInFlightMaxNum(2);
Expand Down
Loading