Skip to content

Commit 06ab4f6

Browse files
authored
[Pipe] Preserve parser fairness and OOM retry state (#18306)
1 parent 2cfccad commit 06ab4f6

4 files changed

Lines changed: 134 additions & 5 deletions

File tree

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/container/scan/TsFileInsertionScanDataContainer.java

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
package org.apache.iotdb.db.pipe.event.common.tsfile.container.scan;
2121

2222
import org.apache.iotdb.commons.exception.IllegalPathException;
23+
import org.apache.iotdb.commons.exception.pipe.PipeRuntimeOutOfMemoryCriticalException;
2324
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
2425
import org.apache.iotdb.commons.pipe.config.PipeConfig;
2526
import org.apache.iotdb.commons.pipe.datastructure.pattern.PipePattern;
@@ -352,6 +353,10 @@ private Tablet getNextTablet() {
352353
}
353354
PipeTabletUtils.compactBitMaps(tablet);
354355
return tablet;
356+
} catch (final PipeRuntimeOutOfMemoryCriticalException e) {
357+
// Keep the parser state so the caller can yield its parser slot and retry from the same
358+
// unconsumed data after memory is available again.
359+
throw e;
355360
} catch (final Exception e) {
356361
close();
357362
throw new PipeException("Failed to get next tablet insertion event.", e);

iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManager.java

Lines changed: 41 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,7 @@ public class PipeMemoryManager {
7171
private final Map<PipeIdentity, ArrayDeque<PipeRegionIdentity>>
7272
waitingTsFileParserRegionOrderByPipe = new HashMap<>();
7373
private final ArrayDeque<PipeIdentity> waitingTsFileParserPipeOrder = new ArrayDeque<>();
74+
private PipeIdentity lastAdmittedWaitingTsFileParserPipe;
7475

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

182184
final int globalLimit = Math.max(1, PIPE_CONFIG.getPipeTsFileParserInFlightMaxNum());
183185
final int perPipeRegionLimit =
@@ -212,6 +214,9 @@ public synchronized boolean tryReserveTsFileParserMemory(
212214
}
213215

214216
removeTsFileParserReservationRequest(pipeRegionIdentity, reservationKey, true);
217+
if (wasRequestAlreadyWaiting) {
218+
lastAdmittedWaitingTsFileParserPipe = pipeIdentity;
219+
}
215220
reservedTsFileParserCount++;
216221
reservedTsFileParserCountByPipe.merge(pipeIdentity, 1, Integer::sum);
217222
reservedTsFileParserCountByPipeRegion.put(pipeRegionIdentity, reservedCountOfPipeRegion + 1);
@@ -262,10 +267,11 @@ public synchronized void releaseTsFileParserMemory(
262267
reservedTsFileParserCountByPipe.put(pipeIdentity, reservedCountOfPipe - 1);
263268
}
264269
reservedTsFileParserCount--;
270+
clearTsFileParserAdmissionCursorIfIdle();
265271
notifyNextTsFileParserMemoryReservationInternal();
266272
}
267273

268-
private void enqueueTsFileParserReservationRequest(
274+
private boolean enqueueTsFileParserReservationRequest(
269275
final PipeRegionIdentity pipeRegionIdentity,
270276
final TsFileParserMemoryReservation reservationKey) {
271277
final LinkedHashSet<TsFileParserMemoryReservation> requestsOfPipeRegion =
@@ -282,7 +288,7 @@ private void enqueueTsFileParserReservationRequest(
282288
regionOrder.addLast(key);
283289
return new LinkedHashSet<>();
284290
});
285-
requestsOfPipeRegion.add(reservationKey);
291+
return !requestsOfPipeRegion.add(reservationKey);
286292
}
287293

288294
public synchronized void notifyNextTsFileParserMemoryReservation() {
@@ -322,7 +328,14 @@ private void notifyNextTsFileParserMemoryReservationInternal() {
322328

323329
private PipeRegionIdentity getNextEligibleTsFileParserPipeRegion(
324330
final int perPipeRegionLimit, final boolean requirePipeWithoutReservedParser) {
331+
PipeRegionIdentity firstEligiblePipeRegion = null;
332+
boolean hasVisitedLastAdmittedPipe = lastAdmittedWaitingTsFileParserPipe == null;
325333
for (final PipeIdentity pipeIdentity : waitingTsFileParserPipeOrder) {
334+
final boolean isLastAdmittedPipe = pipeIdentity.equals(lastAdmittedWaitingTsFileParserPipe);
335+
if (isLastAdmittedPipe) {
336+
hasVisitedLastAdmittedPipe = true;
337+
}
338+
326339
// Under soft memory pressure, reserve the hard-threshold headroom for a pipe that has no
327340
// parser yet. Otherwise a busy pipe at the queue head can block every pipe behind it.
328341
if (requirePipeWithoutReservedParser
@@ -335,14 +348,32 @@ private PipeRegionIdentity getNextEligibleTsFileParserPipeRegion(
335348
if (regionOrder == null) {
336349
continue;
337350
}
351+
PipeRegionIdentity eligiblePipeRegion = null;
338352
for (final PipeRegionIdentity pipeRegionIdentity : regionOrder) {
339353
if (reservedTsFileParserCountByPipeRegion.getOrDefault(pipeRegionIdentity, 0)
340354
< perPipeRegionLimit) {
341-
return pipeRegionIdentity;
355+
eligiblePipeRegion = pipeRegionIdentity;
356+
break;
342357
}
343358
}
359+
if (eligiblePipeRegion == null) {
360+
continue;
361+
}
362+
363+
if (firstEligiblePipeRegion == null) {
364+
firstEligiblePipeRegion = eligiblePipeRegion;
365+
}
366+
if (hasVisitedLastAdmittedPipe && !isLastAdmittedPipe) {
367+
return eligiblePipeRegion;
368+
}
369+
}
370+
return firstEligiblePipeRegion;
371+
}
372+
373+
private void clearTsFileParserAdmissionCursorIfIdle() {
374+
if (reservedTsFileParserCount == 0 && waitingTsFileParserPipeOrder.isEmpty()) {
375+
lastAdmittedWaitingTsFileParserPipe = null;
344376
}
345-
return null;
346377
}
347378

348379
private void removeTsFileParserReservationRequest(
@@ -365,6 +396,9 @@ private void removeTsFileParserReservationRequest(
365396
if (regionOrder.isEmpty()) {
366397
waitingTsFileParserRegionOrderByPipe.remove(pipeIdentity);
367398
waitingTsFileParserPipeOrder.remove(pipeIdentity);
399+
if (!rotateAfterAdmission) {
400+
clearTsFileParserAdmissionCursorIfIdle();
401+
}
368402
return;
369403
}
370404
}
@@ -376,6 +410,8 @@ private void removeTsFileParserReservationRequest(
376410
if (rotateAfterAdmission) {
377411
waitingTsFileParserPipeOrder.remove(pipeIdentity);
378412
waitingTsFileParserPipeOrder.addLast(pipeIdentity);
413+
} else {
414+
clearTsFileParserAdmissionCursorIfIdle();
379415
}
380416
}
381417

iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/TsFileInsertionDataContainerTest.java

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -184,6 +184,47 @@ public void testScanContainerReleasesTabletMemoryAfterRawTabletGenerated() throw
184184
}
185185
}
186186

187+
@Test
188+
public void testScanContainerKeepsIteratorOnOutOfMemory() throws Exception {
189+
nonalignedTsFile =
190+
TsFileGeneratorUtils.generateNonAlignedTsFile(
191+
"nonaligned-retry-tablet-memory.tsfile", 1, 1, 10, 0, 100, 10, 10);
192+
193+
try (final TsFileInsertionScanDataContainer container =
194+
new TsFileInsertionScanDataContainer(
195+
nonalignedTsFile,
196+
new PrefixPipePattern("root"),
197+
Long.MIN_VALUE,
198+
Long.MAX_VALUE,
199+
null,
200+
null,
201+
false)) {
202+
final AtomicInteger memoryUsageReadCount = new AtomicInteger(0);
203+
replaceAllocatedTabletMemory(
204+
container,
205+
new PipeMemoryBlock(0) {
206+
@Override
207+
public long getMemoryUsageInBytes() {
208+
if (memoryUsageReadCount.incrementAndGet() == 2) {
209+
throw new PipeRuntimeOutOfMemoryCriticalException("expected oom");
210+
}
211+
return super.getMemoryUsageInBytes();
212+
}
213+
});
214+
215+
final Iterator<TabletInsertionEvent> iterator =
216+
container.toTabletInsertionEvents().iterator();
217+
final PipeRuntimeOutOfMemoryCriticalException exception =
218+
Assert.assertThrows(PipeRuntimeOutOfMemoryCriticalException.class, iterator::next);
219+
Assert.assertEquals("expected oom", exception.getMessage());
220+
221+
Assert.assertTrue(iterator.hasNext());
222+
final TabletInsertionEvent event = iterator.next();
223+
Assert.assertTrue(event instanceof PipeRawTabletInsertionEvent);
224+
((PipeRawTabletInsertionEvent) event).clearReferenceCount(getClass().getName());
225+
}
226+
}
227+
187228
@Test
188229
public void testConsumeTabletInsertionEventsWithRetryPreservesProgressOnOutOfMemory()
189230
throws Exception {
@@ -1309,6 +1350,16 @@ private PipeMemoryBlock getAllocatedTabletMemory(final TsFileInsertionDataContai
13091350
return (PipeMemoryBlock) field.get(container);
13101351
}
13111352

1353+
private void replaceAllocatedTabletMemory(
1354+
final TsFileInsertionDataContainer container, final PipeMemoryBlock replacement)
1355+
throws NoSuchFieldException, IllegalAccessException {
1356+
final Field field =
1357+
TsFileInsertionDataContainer.class.getDeclaredField("allocatedMemoryBlockForTablet");
1358+
field.setAccessible(true);
1359+
((PipeMemoryBlock) field.get(container)).close();
1360+
field.set(container, replacement);
1361+
}
1362+
13121363
@SuppressWarnings("unchecked")
13131364
private AtomicReference<TsFileInsertionDataContainer> getDataContainer(
13141365
final PipeTsFileInsertionEvent event) throws NoSuchFieldException, IllegalAccessException {

iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/resource/memory/PipeMemoryManagerTest.java

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -192,6 +192,43 @@ public void testPipeFairnessIsNotWeightedByRegionCount() {
192192
Assert.assertTrue(tryAcquire(pipeARegion2));
193193
}
194194

195+
@Test
196+
public void testPipeFairnessSurvivesTransientSingleRegionQueueGap() {
197+
commonConfig.setPipeTsFileParserInFlightMaxNum(1);
198+
commonConfig.setPipeTsFileParserInFlightMaxNumPerPipeRegion(1);
199+
200+
final Reservation blocker = new Reservation("blocker", 0);
201+
final Reservation multiRegion1 = new Reservation("multi", 1, "1");
202+
final Reservation multiRegion2 = new Reservation("multi", 1, "2");
203+
final Reservation multiRegion3 = new Reservation("multi", 1, "3");
204+
final Reservation singleFirst = new Reservation("single", 2, "1");
205+
final Reservation singleSecond = new Reservation("single", 2, "1");
206+
207+
Assert.assertTrue(tryAcquire(blocker));
208+
Assert.assertFalse(tryAcquire(multiRegion1));
209+
Assert.assertFalse(tryAcquire(multiRegion2));
210+
Assert.assertFalse(tryAcquire(multiRegion3));
211+
Assert.assertFalse(tryAcquire(singleFirst));
212+
213+
release(blocker);
214+
Assert.assertTrue(tryAcquire(multiRegion1));
215+
release(multiRegion1);
216+
217+
Assert.assertTrue(tryAcquire(singleFirst));
218+
release(singleFirst);
219+
220+
// The single-region pipe temporarily has no admission request while it advances to its next
221+
// TsFile, so the multi-region pipe can use the otherwise idle parser slot.
222+
Assert.assertTrue(tryAcquire(multiRegion2));
223+
Assert.assertFalse(tryAcquire(singleSecond));
224+
release(multiRegion2);
225+
226+
// Once the single-region pipe is waiting again, the remembered pipe-level cursor must prevent
227+
// another region of the multi-region pipe from taking a second consecutive turn.
228+
Assert.assertFalse(tryAcquire(multiRegion3));
229+
Assert.assertTrue(tryAcquire(singleSecond));
230+
}
231+
195232
@Test
196233
public void testSoftMemoryHeadroomIsReservedForPipeWithoutParser() {
197234
commonConfig.setPipeTsFileParserInFlightMaxNum(2);

0 commit comments

Comments
 (0)