Skip to content
Closed
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 @@ -33,6 +33,7 @@
import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryBlock;
import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryWeightUtil;
import org.apache.iotdb.db.storageengine.dataregion.modification.ModEntry;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import org.apache.iotdb.db.utils.datastructure.PatternTreeMapFactory;
import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;

Expand All @@ -45,6 +46,7 @@

import java.io.File;
import java.io.IOException;
import java.util.List;

public abstract class TsFileInsertionEventParser implements AutoCloseable {

Expand Down Expand Up @@ -78,6 +80,8 @@ public abstract class TsFileInsertionEventParser implements AutoCloseable {

protected Iterable<TabletInsertionEvent> tabletInsertionIterable;

protected final boolean objectPathsOnly;

protected TsFileInsertionEventParser(
final File tsFile,
final String pipeName,
Expand All @@ -91,6 +95,38 @@ protected TsFileInsertionEventParser(
final boolean skipIfNoPrivileges,
final PipeInsertionEvent sourceEvent,
final boolean isWithMod) {
this(
tsFile,
pipeName,
creationTime,
treePattern,
tablePattern,
startTime,
endTime,
pipeTaskMeta,
entity,
skipIfNoPrivileges,
sourceEvent,
null,
false,
isWithMod);
}

protected TsFileInsertionEventParser(
final File tsFile,
final String pipeName,
final long creationTime,
final TreePattern treePattern,
final TablePattern tablePattern,
final long startTime,
final long endTime,
final PipeTaskMeta pipeTaskMeta,
final IAuditEntity entity,
final boolean skipIfNoPrivileges,
final PipeInsertionEvent sourceEvent,
final TsFileResource tsFileResource,
final boolean objectPathsOnly,
final boolean isWithMod) {
this.pipeName = pipeName;
this.creationTime = creationTime;
this.entity = entity;
Expand All @@ -107,6 +143,7 @@ protected TsFileInsertionEventParser(

this.pipeTaskMeta = pipeTaskMeta;
this.sourceEvent = sourceEvent;
this.objectPathsOnly = objectPathsOnly;

this.allocatedMemoryBlockForTablet =
PipeDataNodeResourceManager.memory()
Expand All @@ -130,6 +167,8 @@ protected TsFileInsertionEventParser(
*/
public abstract Iterable<TabletInsertionEvent> toTabletInsertionEvents();

public void drainGeneratedObjectColumnModEntriesTo(final List<ModEntry> modEntries) {}

/**
* Record parse start time when hasNext() is called for the first time and returns true. Should be
* called in Iterator.hasNext() when it's the first call.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ public TsFileInsertionEventQueryParser(
final long endTime,
final PipeInsertionEvent sourceEvent)
throws IOException, IllegalPathException {
this(null, 0, tsFile, pattern, startTime, endTime, null, sourceEvent, false);
this(null, 0, tsFile, pattern, startTime, endTime, null, sourceEvent, false, false);
}

public TsFileInsertionEventQueryParser(
Expand All @@ -104,6 +104,31 @@ public TsFileInsertionEventQueryParser(
final PipeInsertionEvent sourceEvent,
final boolean isWithMod)
throws IOException, IllegalPathException {
this(
pipeName,
creationTime,
tsFile,
pattern,
startTime,
endTime,
pipeTaskMeta,
sourceEvent,
isWithMod,
false);
}

public TsFileInsertionEventQueryParser(
final String pipeName,
final long creationTime,
final File tsFile,
final TreePattern pattern,
final long startTime,
final long endTime,
final PipeTaskMeta pipeTaskMeta,
final PipeInsertionEvent sourceEvent,
final boolean isWithMod,
final boolean objectPathsOnly)
throws IOException, IllegalPathException {
this(
pipeName,
creationTime,
Expand All @@ -116,7 +141,8 @@ public TsFileInsertionEventQueryParser(
null,
false,
null,
isWithMod);
isWithMod,
objectPathsOnly);
}

public TsFileInsertionEventQueryParser(
Expand All @@ -133,6 +159,37 @@ public TsFileInsertionEventQueryParser(
final Map<IDeviceID, Boolean> deviceIsAlignedMap,
final boolean isWithMod)
throws IOException, IllegalPathException {
this(
pipeName,
creationTime,
tsFile,
pattern,
startTime,
endTime,
pipeTaskMeta,
sourceEvent,
entity,
skipIfNoPrivileges,
deviceIsAlignedMap,
isWithMod,
false);
}

public TsFileInsertionEventQueryParser(
final String pipeName,
final long creationTime,
final File tsFile,
final TreePattern pattern,
final long startTime,
final long endTime,
final PipeTaskMeta pipeTaskMeta,
final PipeInsertionEvent sourceEvent,
final IAuditEntity entity,
final boolean skipIfNoPrivileges,
final Map<IDeviceID, Boolean> deviceIsAlignedMap,
final boolean isWithMod,
final boolean objectPathsOnly)
throws IOException, IllegalPathException {
super(
tsFile,
pipeName,
Expand All @@ -145,6 +202,8 @@ public TsFileInsertionEventQueryParser(
entity,
skipIfNoPrivileges,
sourceEvent,
null,
objectPathsOnly,
isWithMod);

try {
Expand All @@ -157,7 +216,7 @@ public TsFileInsertionEventQueryParser(
.forceAllocateForTabletWithRetry(currentModifications.ramBytesUsed());

final PipeTsFileResourceManager tsFileResourceManager = PipeDataNodeResourceManager.tsfile();
final Map<IDeviceID, List<String>> deviceMeasurementsMap;
Map<IDeviceID, List<String>> deviceMeasurementsMap;

// TsFileReader is not thread-safe, so we need to create it here and close it later.
long memoryRequiredInBytes =
Expand Down Expand Up @@ -197,6 +256,10 @@ public TsFileInsertionEventQueryParser(
memoryRequiredInBytes +=
PipeMemoryWeightUtil.memoryOfIDeviceID2StrList(deviceMeasurementsMap);
}
if (objectPathsOnly) {
deviceMeasurementsMap =
filterObjectMeasurementsMap(deviceMeasurementsMap, measurementDataTypeMap);
}
allocatedMemoryBlock =
PipeDataNodeResourceManager.memory().forceAllocate(memoryRequiredInBytes);

Expand Down Expand Up @@ -373,9 +436,10 @@ private Map<IDeviceID, List<String>> readFilteredDeviceMeasurementsMap(
final Map<IDeviceID, List<String>> result = new HashMap<>();

for (final IDeviceID device : devices) {
tsFileSequenceReader
.readDeviceMetadata(device)
.values()
tsFileSequenceReader.readDeviceMetadata(device).values().stream()
.filter(
timeseriesMetadata ->
!objectPathsOnly || timeseriesMetadata.getTsDataType() == TSDataType.OBJECT)
.forEach(
timeseriesMetadata ->
result
Expand All @@ -386,6 +450,24 @@ private Map<IDeviceID, List<String>> readFilteredDeviceMeasurementsMap(
return result;
}

private Map<IDeviceID, List<String>> filterObjectMeasurementsMap(
final Map<IDeviceID, List<String>> deviceMeasurementsMap,
final Map<String, TSDataType> measurementDataTypeMap) {
final Map<IDeviceID, List<String>> result = new HashMap<>();
for (final Map.Entry<IDeviceID, List<String>> entry : deviceMeasurementsMap.entrySet()) {
final List<String> measurements = new ArrayList<>();
for (final String measurement : entry.getValue()) {
if (measurementDataTypeMap.get(entry.getKey() + "." + measurement) == TSDataType.OBJECT) {
measurements.add(measurement);
}
}
if (!measurements.isEmpty()) {
result.put(entry.getKey(), measurements);
}
}
return result;
}

@Override
public Iterable<TabletInsertionEvent> toTabletInsertionEvents() {
if (tabletInsertionIterable == null) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,35 @@ public TsFileInsertionEventScanParser(
final PipeInsertionEvent sourceEvent,
final boolean isWithMod)
throws IOException, IllegalPathException {
this(
pipeName,
creationTime,
tsFile,
pattern,
startTime,
endTime,
pipeTaskMeta,
entity,
skipIfNoPrivileges,
sourceEvent,
isWithMod,
false);
}

public TsFileInsertionEventScanParser(
final String pipeName,
final long creationTime,
final File tsFile,
final TreePattern pattern,
final long startTime,
final long endTime,
final PipeTaskMeta pipeTaskMeta,
final IAuditEntity entity,
final boolean skipIfNoPrivileges,
final PipeInsertionEvent sourceEvent,
final boolean isWithMod,
final boolean objectPathsOnly)
throws IOException, IllegalPathException {
super(
tsFile,
pipeName,
Expand All @@ -135,6 +164,8 @@ public TsFileInsertionEventScanParser(
entity,
skipIfNoPrivileges,
sourceEvent,
null,
objectPathsOnly,
isWithMod);

this.startTime = startTime;
Expand Down Expand Up @@ -181,6 +212,19 @@ public TsFileInsertionEventScanParser(
final PipeInsertionEvent sourceEvent,
final boolean isWithMod)
throws IOException, IllegalPathException {
this(tsFile, pattern, startTime, endTime, pipeTaskMeta, sourceEvent, isWithMod, false);
}

public TsFileInsertionEventScanParser(
final File tsFile,
final TreePattern pattern,
final long startTime,
final long endTime,
final PipeTaskMeta pipeTaskMeta,
final PipeInsertionEvent sourceEvent,
final boolean isWithMod,
final boolean objectPathsOnly)
throws IOException, IllegalPathException {
this(
null,
0,
Expand All @@ -192,7 +236,8 @@ public TsFileInsertionEventScanParser(
null,
false,
sourceEvent,
isWithMod);
isWithMod,
objectPathsOnly);
}

@Override
Expand Down Expand Up @@ -416,8 +461,12 @@ private boolean putValueToColumns(final BatchData data, final Tablet tablet, fin
case TEXT:
case BLOB:
case STRING:
case OBJECT:
PipeTabletUtils.putValue(
tablet, rowIndex, i, tablet.getSchemas().get(i).getType(), Binary.EMPTY_VALUE);
break;
default:
break;
}
PipeTabletUtils.markNullValue(tablet, rowIndex, i);
continue;
Expand Down Expand Up @@ -469,6 +518,7 @@ private boolean putValueToColumns(final BatchData data, final Tablet tablet, fin
case TEXT:
case BLOB:
case STRING:
case OBJECT:
final Binary binary = primitiveType.getBinary();
PipeTabletUtils.putValue(
tablet,
Expand Down Expand Up @@ -524,6 +574,7 @@ private boolean putValueToColumns(final BatchData data, final Tablet tablet, fin
case TEXT:
case BLOB:
case STRING:
case OBJECT:
final Binary binary = data.getBinary();
PipeTabletUtils.putValue(
tablet,
Expand Down Expand Up @@ -773,6 +824,11 @@ private boolean filterChunk(
return true;
}

if (objectPathsOnly && chunkHeader.getDataType() != TSDataType.OBJECT) {
tsFileSequenceReader.position(nextMarkerOffset);
return true;
}

// Skip the chunk if it is fully deleted by mods
if (!currentModifications.isEmpty()) {
final Statistics statistics =
Expand Down
Loading
Loading