Skip to content

Commit a7c7e54

Browse files
authored
[To dev/1.3] Load: Optimize fallback logic for corrupted TsFile (#18004)
* Load: Optimized the downgraded logic for tsFile to insert more data when tsFile corrupted (#17674) * down * Update LoadTreeStatementDataTypeConvertExecutionVisitorTest.java * Fix query parser deferred tablet failure handling * Fix aligned load fallback corruption test * Address load fallback review comments (cherry picked from commit cb97fe4) * Update LoadTreeStatementDataTypeConvertExecutionVisitorTest.java * Update LoadTreeStatementDataTypeConvertExecutionVisitorTest.java
1 parent 2717029 commit a7c7e54

5 files changed

Lines changed: 1106 additions & 28 deletions

File tree

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/container/query/TsFileInsertionQueryDataContainer.java‎

Lines changed: 95 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -51,9 +51,9 @@
5151
import java.io.File;
5252
import java.io.IOException;
5353
import java.util.ArrayList;
54-
import java.util.HashMap;
5554
import java.util.HashSet;
5655
import java.util.Iterator;
56+
import java.util.LinkedHashMap;
5757
import java.util.List;
5858
import java.util.Map;
5959
import java.util.NoSuchElementException;
@@ -72,6 +72,7 @@ public class TsFileInsertionQueryDataContainer extends TsFileInsertionDataContai
7272
private final Map<IDeviceID, Boolean> deviceIsAlignedMap;
7373
private final Map<String, TSDataType> measurementDataTypeMap;
7474
private final TabletStringInternPool tabletStringInternPool = new TabletStringInternPool();
75+
private RuntimeException deferredException;
7576

7677
@TestOnly
7778
public TsFileInsertionQueryDataContainer(
@@ -101,6 +102,7 @@ public TsFileInsertionQueryDataContainer(
101102
pipeTaskMeta,
102103
sourceEvent,
103104
null,
105+
null,
104106
isWithMod);
105107
}
106108

@@ -116,6 +118,33 @@ public TsFileInsertionQueryDataContainer(
116118
final Map<IDeviceID, Boolean> deviceIsAlignedMap,
117119
final boolean isWithMod)
118120
throws IOException {
121+
this(
122+
pipeName,
123+
creationTime,
124+
tsFile,
125+
pattern,
126+
startTime,
127+
endTime,
128+
pipeTaskMeta,
129+
sourceEvent,
130+
deviceIsAlignedMap,
131+
null,
132+
isWithMod);
133+
}
134+
135+
public TsFileInsertionQueryDataContainer(
136+
final String pipeName,
137+
final long creationTime,
138+
final File tsFile,
139+
final PipePattern pattern,
140+
final long startTime,
141+
final long endTime,
142+
final PipeTaskMeta pipeTaskMeta,
143+
final EnrichedEvent sourceEvent,
144+
final Map<IDeviceID, Boolean> deviceIsAlignedMap,
145+
final Map<IDeviceID, List<String>> deviceMeasurementsMapOverride,
146+
final boolean isWithMod)
147+
throws IOException {
119148
super(
120149
tsFile,
121150
pipeName,
@@ -145,7 +174,25 @@ public TsFileInsertionQueryDataContainer(
145174
tsFileSequenceReader = new TsFileSequenceReader(tsFile.getPath(), true, true);
146175
tsFileReader = new TsFileReader(tsFileSequenceReader);
147176

148-
if (tsFileResourceManager.cacheObjectsIfAbsent(tsFile)) {
177+
if (Objects.nonNull(deviceMeasurementsMapOverride)) {
178+
this.deviceIsAlignedMap =
179+
Objects.nonNull(deviceIsAlignedMap)
180+
? new LinkedHashMap<>(deviceIsAlignedMap)
181+
: readDeviceIsAlignedMap();
182+
memoryRequiredInBytes +=
183+
Objects.nonNull(deviceIsAlignedMap)
184+
? 0
185+
: PipeMemoryWeightUtil.memoryOfIDeviceId2Bool(this.deviceIsAlignedMap);
186+
187+
measurementDataTypeMap =
188+
readFilteredFullPathDataTypeMap(deviceMeasurementsMapOverride.keySet());
189+
memoryRequiredInBytes +=
190+
PipeMemoryWeightUtil.memoryOfStr2TSDataType(measurementDataTypeMap);
191+
192+
deviceMeasurementsMap = new LinkedHashMap<>(deviceMeasurementsMapOverride);
193+
memoryRequiredInBytes +=
194+
PipeMemoryWeightUtil.memoryOfIDeviceID2StrList(deviceMeasurementsMap);
195+
} else if (tsFileResourceManager.cacheObjectsIfAbsent(tsFile)) {
149196
// These read-only objects can be found in cache.
150197
this.deviceIsAlignedMap =
151198
Objects.nonNull(deviceIsAlignedMap)
@@ -246,9 +293,31 @@ public TsFileInsertionQueryDataContainer(
246293
}
247294
}
248295

296+
public TsFileInsertionQueryDataContainer(
297+
final File tsFile,
298+
final PipePattern pattern,
299+
final long startTime,
300+
final long endTime,
301+
final Map<IDeviceID, List<String>> deviceMeasurementsMapOverride,
302+
final boolean isWithMod)
303+
throws IOException {
304+
this(
305+
null,
306+
0,
307+
tsFile,
308+
pattern,
309+
startTime,
310+
endTime,
311+
null,
312+
null,
313+
null,
314+
deviceMeasurementsMapOverride,
315+
isWithMod);
316+
}
317+
249318
private Map<IDeviceID, List<String>> filterDeviceMeasurementsMapByPattern(
250319
final Map<IDeviceID, List<String>> originalDeviceMeasurementsMap) {
251-
final Map<IDeviceID, List<String>> filteredDeviceMeasurementsMap = new HashMap<>();
320+
final Map<IDeviceID, List<String>> filteredDeviceMeasurementsMap = new LinkedHashMap<>();
252321
for (final Map.Entry<IDeviceID, List<String>> entry :
253322
originalDeviceMeasurementsMap.entrySet()) {
254323
final String deviceId = ((PlainDeviceID) entry.getKey()).toStringID();
@@ -282,7 +351,7 @@ else if (pattern.mayOverlapWithDevice(deviceId)) {
282351
}
283352

284353
private Map<IDeviceID, Boolean> readDeviceIsAlignedMap() throws IOException {
285-
final Map<IDeviceID, Boolean> deviceIsAlignedResultMap = new HashMap<>();
354+
final Map<IDeviceID, Boolean> deviceIsAlignedResultMap = new LinkedHashMap<>();
286355
final TsFileDeviceIterator deviceIsAlignedIterator =
287356
tsFileSequenceReader.getAllDevicesIteratorWithIsAligned();
288357
while (deviceIsAlignedIterator.hasNext()) {
@@ -313,7 +382,7 @@ private Set<IDeviceID> filterDevicesByPattern(final Set<IDeviceID> devices) {
313382
*/
314383
private Map<String, TSDataType> readFilteredFullPathDataTypeMap(final Set<IDeviceID> devices)
315384
throws IOException {
316-
final Map<String, TSDataType> result = new HashMap<>();
385+
final Map<String, TSDataType> result = new LinkedHashMap<>();
317386

318387
for (final IDeviceID device : devices) {
319388
tsFileSequenceReader
@@ -337,7 +406,7 @@ private Map<String, TSDataType> readFilteredFullPathDataTypeMap(final Set<IDevic
337406
*/
338407
private Map<IDeviceID, List<String>> readFilteredDeviceMeasurementsMap(
339408
final Set<IDeviceID> devices) throws IOException {
340-
final Map<IDeviceID, List<String>> result = new HashMap<>();
409+
final Map<IDeviceID, List<String>> result = new LinkedHashMap<>();
341410

342411
for (final IDeviceID device : devices) {
343412
tsFileSequenceReader
@@ -364,6 +433,7 @@ public Iterable<TabletInsertionEvent> toTabletInsertionEvents() {
364433

365434
@Override
366435
public boolean hasNext() {
436+
throwIfDeferredException();
367437
boolean hasNext = false;
368438
while (tabletIterator == null || !tabletIterator.hasNext()) {
369439
if (!deviceMeasurementsMapIterator.hasNext()) {
@@ -416,9 +486,16 @@ public TabletInsertionEvent next() {
416486
recordTabletMetrics(tablet);
417487
final boolean isAligned =
418488
deviceIsAlignedMap.getOrDefault(new PlainDeviceID(tablet.deviceId), false);
489+
boolean isLast;
490+
try {
491+
isLast = !hasNext();
492+
} catch (final RuntimeException e) {
493+
deferredException = e;
494+
isLast = false;
495+
}
419496

420497
final TabletInsertionEvent next;
421-
if (!hasNext()) {
498+
if (isLast) {
422499
next =
423500
new PipeRawTabletInsertionEvent(
424501
tablet,
@@ -448,8 +525,19 @@ public TabletInsertionEvent next() {
448525
return tabletInsertionIterable;
449526
}
450527

528+
private void throwIfDeferredException() {
529+
if (Objects.isNull(deferredException)) {
530+
return;
531+
}
532+
533+
final RuntimeException exception = deferredException;
534+
deferredException = null;
535+
throw exception;
536+
}
537+
451538
@Override
452539
public void close() {
540+
deferredException = null;
453541
try {
454542
if (tsFileReader != null) {
455543
tsFileReader.close();

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

Lines changed: 52 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -45,6 +45,7 @@
4545
import org.apache.tsfile.file.header.ChunkHeader;
4646
import org.apache.tsfile.file.metadata.AlignedChunkMetadata;
4747
import org.apache.tsfile.file.metadata.IChunkMetadata;
48+
import org.apache.tsfile.file.metadata.IDeviceID;
4849
import org.apache.tsfile.file.metadata.PlainDeviceID;
4950
import org.apache.tsfile.file.metadata.statistics.Statistics;
5051
import org.apache.tsfile.read.TsFileSequenceReader;
@@ -96,6 +97,7 @@ public class TsFileInsertionScanDataContainer extends TsFileInsertionDataContain
9697
private boolean currentIsAligned;
9798
private final List<MeasurementSchema> currentMeasurements = new ArrayList<>();
9899
private final TabletStringInternPool tabletStringInternPool = new TabletStringInternPool();
100+
private Exception deferredException;
99101

100102
private final List<ModsOperationUtil.ModsInfo> modsInfos = new ArrayList<>();
101103
// Cached time chunk
@@ -187,6 +189,7 @@ public Iterable<TabletInsertionEvent> toTabletInsertionEvents() {
187189

188190
@Override
189191
public boolean hasNext() {
192+
throwIfDeferredException();
190193
final boolean hasNext = Objects.nonNull(chunkReader);
191194
if (hasNext && !parseStartTimeRecorded) {
192195
// Record start time on first hasNext() that returns true
@@ -205,7 +208,7 @@ public TabletInsertionEvent next() {
205208
throw new NoSuchElementException();
206209
}
207210

208-
// currentIsAligned is initialized when TsFileInsertionEventScanParser is
211+
// currentIsAligned is initialized when TsFileInsertionScanDataContainer is
209212
// constructed.
210213
// When the getNextTablet function is called, currentIsAligned may be updated,
211214
// causing
@@ -215,7 +218,7 @@ public TabletInsertionEvent next() {
215218
final Tablet tablet = getNextTablet();
216219
// Record tablet metrics
217220
recordTabletMetrics(tablet);
218-
final boolean hasNext = hasNext();
221+
final boolean isLast = isLastTabletWithoutDeferredException();
219222
try {
220223
return new PipeRawTabletInsertionEvent(
221224
tablet,
@@ -224,9 +227,10 @@ public TabletInsertionEvent next() {
224227
sourceEvent != null ? sourceEvent.getCreationTime() : 0,
225228
pipeTaskMeta,
226229
sourceEvent,
227-
!hasNext);
230+
isLast);
228231
} finally {
229-
if (!hasNext) {
232+
if (isLast) {
233+
recordParseEndTimeIfNecessary();
230234
close();
231235
}
232236
}
@@ -241,6 +245,7 @@ public Iterable<Pair<Tablet, Boolean>> toTabletWithIsAligneds() {
241245
new Iterator<Pair<Tablet, Boolean>>() {
242246
@Override
243247
public boolean hasNext() {
248+
throwIfDeferredException();
244249
return Objects.nonNull(chunkReader);
245250
}
246251

@@ -251,24 +256,39 @@ public Pair<Tablet, Boolean> next() {
251256
throw new NoSuchElementException();
252257
}
253258

254-
// currentIsAligned is initialized when TsFileInsertionEventScanParser is constructed.
259+
// currentIsAligned is initialized when TsFileInsertionScanDataContainer is constructed.
255260
// When the getNextTablet function is called, currentIsAligned may be updated, causing
256261
// the currentIsAligned information to be inconsistent with the current Tablet
257262
// information.
258263
final boolean isAligned = currentIsAligned;
259264
final Tablet tablet = getNextTablet();
260-
final boolean hasNext = hasNext();
261265
try {
262266
return new Pair<>(tablet, isAligned);
263267
} finally {
264-
if (!hasNext) {
268+
if (isLastTabletWithoutDeferredException()) {
265269
close();
266270
}
267271
}
268272
}
269273
};
270274
}
271275

276+
public IDeviceID getCurrentDevice() {
277+
return Objects.nonNull(currentDevice) ? new PlainDeviceID(currentDevice) : null;
278+
}
279+
280+
public boolean isCurrentAligned() {
281+
return currentIsAligned;
282+
}
283+
284+
public List<String> getCurrentMeasurements() {
285+
final List<String> measurementIds = new ArrayList<>(currentMeasurements.size());
286+
for (final MeasurementSchema schema : currentMeasurements) {
287+
measurementIds.add(schema.getMeasurementId());
288+
}
289+
return measurementIds;
290+
}
291+
272292
private Tablet getNextTablet() {
273293
try {
274294
Tablet tablet = null;
@@ -321,7 +341,11 @@ private Tablet getNextTablet() {
321341

322342
// Switch chunk reader iff current chunk is all consumed
323343
if (!data.hasCurrent()) {
324-
prepareData();
344+
try {
345+
prepareData();
346+
} catch (final Exception e) {
347+
deferredException = e;
348+
}
325349
}
326350
PipeTabletUtils.compactBitMaps(tablet);
327351
return tablet;
@@ -331,6 +355,26 @@ private Tablet getNextTablet() {
331355
}
332356
}
333357

358+
private void throwIfDeferredException() {
359+
if (Objects.isNull(deferredException)) {
360+
return;
361+
}
362+
363+
final Exception exception = deferredException;
364+
deferredException = null;
365+
throw new PipeException("Failed to prepare next tablet insertion event.", exception);
366+
}
367+
368+
private boolean isLastTabletWithoutDeferredException() {
369+
return Objects.isNull(deferredException) && Objects.isNull(chunkReader);
370+
}
371+
372+
private void recordParseEndTimeIfNecessary() {
373+
if (parseStartTimeRecorded && !parseEndTimeRecorded) {
374+
recordParseEndTime();
375+
}
376+
}
377+
334378
private void prepareData() throws IOException {
335379
do {
336380
do {

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java‎

Lines changed: 3 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -20,9 +20,7 @@
2020
package org.apache.iotdb.db.storageengine.load.converter;
2121

2222
import org.apache.iotdb.common.rpc.thrift.TSStatus;
23-
import org.apache.iotdb.commons.pipe.datastructure.pattern.IoTDBPipePattern;
2423
import org.apache.iotdb.db.conf.IoTDBDescriptor;
25-
import org.apache.iotdb.db.pipe.event.common.tsfile.container.scan.TsFileInsertionScanDataContainer;
2624
import org.apache.iotdb.db.pipe.sink.payload.evolvable.request.PipeTransferTabletRawReq;
2725
import org.apache.iotdb.db.queryengine.plan.statement.Statement;
2826
import org.apache.iotdb.db.queryengine.plan.statement.StatementNode;
@@ -90,17 +88,9 @@ public Optional<TSStatus> visitLoadFile(
9088

9189
try {
9290
for (final File file : loadTsFileStatement.getTsFiles()) {
93-
try (final TsFileInsertionScanDataContainer container =
94-
new TsFileInsertionScanDataContainer(
95-
file,
96-
new IoTDBPipePattern(null),
97-
Long.MIN_VALUE,
98-
Long.MAX_VALUE,
99-
null,
100-
null,
101-
true)) {
102-
for (final Pair<Tablet, Boolean> tabletWithIsAligned :
103-
container.toTabletWithIsAligneds()) {
91+
try (final LoadTreeTsFileTabletIterator tabletIterator =
92+
new LoadTreeTsFileTabletIterator(file, true)) {
93+
for (final Pair<Tablet, Boolean> tabletWithIsAligned : tabletIterator) {
10494
final PipeTransferTabletRawReq tabletRawReq =
10595
PipeTransferTabletRawReq.toTPipeTransferRawReq(
10696
tabletWithIsAligned.getLeft(), tabletWithIsAligned.getRight());

0 commit comments

Comments
 (0)