This is an automated email from the ASF dual-hosted git repository.
rong pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 332521a3241 [IOTDB-6083] Pipe: Fix subscrption running with the
pattern option causing OOM & make PipeRawTabletInsertionEvent able to report
progress to avoid losing data (#10865)
332521a3241 is described below
commit 332521a3241f17e2b10c75d779266d7d31c2f77e
Author: 马子坤 <[email protected]>
AuthorDate: Wed Aug 23 17:52:01 2023 +0800
[IOTDB-6083] Pipe: Fix subscrption running with the pattern option causing
OOM & make PipeRawTabletInsertionEvent able to report progress to avoid losing
data (#10865)
Co-authored-by: Steve Yurong Su <[email protected]>
---
.../protocol/airgap/IoTDBAirGapConnector.java | 35 +++++--
.../thrift/async/IoTDBThriftAsyncConnector.java | 17 ++++
.../thrift/sync/IoTDBThriftSyncConnector.java | 18 ++++
.../apache/iotdb/db/pipe/event/EnrichedEvent.java | 14 ++-
.../db/pipe/event/common/row/PipeRowCollector.java | 18 +++-
.../tablet/PipeInsertNodeTabletInsertionEvent.java | 15 ++-
.../common/tablet/PipeRawTabletInsertionEvent.java | 102 +++++++++++++++++----
.../tablet/TabletInsertionDataContainer.java | 27 +++++-
.../common/tsfile/PipeTsFileInsertionEvent.java | 4 +-
.../tsfile/TsFileInsertionDataContainer.java | 34 ++++++-
.../db/pipe/processor/PipeDoNothingProcessor.java | 49 +---------
.../pipe/task/connection/PipeEventCollector.java | 31 +------
12 files changed, 250 insertions(+), 114 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/airgap/IoTDBAirGapConnector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/airgap/IoTDBAirGapConnector.java
index 9830e44c9e7..40968f4fa9f 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/airgap/IoTDBAirGapConnector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/airgap/IoTDBAirGapConnector.java
@@ -28,6 +28,7 @@ import
org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransfer
import
org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferInsertNodeReq;
import
org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferTabletReq;
import org.apache.iotdb.db.pipe.connector.protocol.IoTDBConnector;
+import org.apache.iotdb.db.pipe.event.EnrichedEvent;
import
org.apache.iotdb.db.pipe.event.common.tablet.PipeInsertNodeTabletInsertionEvent;
import
org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
@@ -173,6 +174,25 @@ public class IoTDBAirGapConnector extends IoTDBConnector {
@Override
public void transfer(TabletInsertionEvent tabletInsertionEvent) throws
Exception {
// PipeProcessor can change the type of TabletInsertionEvent
+ if (!(tabletInsertionEvent instanceof PipeInsertNodeTabletInsertionEvent)
+ && !(tabletInsertionEvent instanceof PipeRawTabletInsertionEvent)) {
+ LOGGER.warn(
+ "IoTDBAirGapConnector only support "
+ + "PipeInsertNodeTabletInsertionEvent and
PipeRawTabletInsertionEvent. "
+ + "Ignore {}.",
+ tabletInsertionEvent);
+ return;
+ }
+
+ if (((EnrichedEvent) tabletInsertionEvent).shouldParsePattern()) {
+ if (tabletInsertionEvent instanceof PipeInsertNodeTabletInsertionEvent) {
+ transfer(
+ ((PipeInsertNodeTabletInsertionEvent)
tabletInsertionEvent).parseEventWithPattern());
+ } else { // tabletInsertionEvent instanceof PipeRawTabletInsertionEvent
+ transfer(((PipeRawTabletInsertionEvent)
tabletInsertionEvent).parseEventWithPattern());
+ }
+ return;
+ }
final int socketIndex = nextSocketIndex();
final Socket socket = sockets.get(socketIndex);
@@ -180,14 +200,8 @@ public class IoTDBAirGapConnector extends IoTDBConnector {
try {
if (tabletInsertionEvent instanceof PipeInsertNodeTabletInsertionEvent) {
doTransfer(socket, (PipeInsertNodeTabletInsertionEvent)
tabletInsertionEvent);
- } else if (tabletInsertionEvent instanceof PipeRawTabletInsertionEvent) {
- doTransfer(socket, (PipeRawTabletInsertionEvent) tabletInsertionEvent);
} else {
- LOGGER.warn(
- "IoTDBAirGapConnector only support "
- + "PipeInsertNodeTabletInsertionEvent and
PipeRawTabletInsertionEvent. "
- + "Ignore {}.",
- tabletInsertionEvent);
+ doTransfer(socket, (PipeRawTabletInsertionEvent) tabletInsertionEvent);
}
} catch (IOException e) {
isSocketAlive.set(socketIndex, false);
@@ -210,6 +224,13 @@ public class IoTDBAirGapConnector extends IoTDBConnector {
return;
}
+ if (((EnrichedEvent) tsFileInsertionEvent).shouldParsePattern()) {
+ for (final TabletInsertionEvent event :
tsFileInsertionEvent.toTabletInsertionEvents()) {
+ transfer(event);
+ }
+ return;
+ }
+
final int socketIndex = nextSocketIndex();
final Socket socket = sockets.get(socketIndex);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java
index b588a0b168d..c649485fa86 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/async/IoTDBThriftAsyncConnector.java
@@ -157,6 +157,16 @@ public class IoTDBThriftAsyncConnector extends
IoTDBConnector {
return;
}
+ if (((EnrichedEvent) tabletInsertionEvent).shouldParsePattern()) {
+ if (tabletInsertionEvent instanceof PipeInsertNodeTabletInsertionEvent) {
+ transfer(
+ ((PipeInsertNodeTabletInsertionEvent)
tabletInsertionEvent).parseEventWithPattern());
+ } else { // tabletInsertionEvent instanceof PipeRawTabletInsertionEvent
+ transfer(((PipeRawTabletInsertionEvent)
tabletInsertionEvent).parseEventWithPattern());
+ }
+ return;
+ }
+
final long requestCommitId = commitIdGenerator.incrementAndGet();
if (isTabletBatchModeEnabled) {
@@ -290,6 +300,13 @@ public class IoTDBThriftAsyncConnector extends
IoTDBConnector {
return;
}
+ if (((EnrichedEvent) tsFileInsertionEvent).shouldParsePattern()) {
+ for (final TabletInsertionEvent event :
tsFileInsertionEvent.toTabletInsertionEvents()) {
+ transfer(event);
+ }
+ return;
+ }
+
final long requestCommitId = commitIdGenerator.incrementAndGet();
final PipeTsFileInsertionEvent pipeTsFileInsertionEvent =
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java
index c4764219859..5a1e4c88db1 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/connector/protocol/thrift/sync/IoTDBThriftSyncConnector.java
@@ -31,6 +31,7 @@ import
org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransfer
import
org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferInsertNodeReq;
import
org.apache.iotdb.db.pipe.connector.payload.evolvable.request.PipeTransferTabletReq;
import org.apache.iotdb.db.pipe.connector.protocol.IoTDBConnector;
+import org.apache.iotdb.db.pipe.event.EnrichedEvent;
import
org.apache.iotdb.db.pipe.event.common.tablet.PipeInsertNodeTabletInsertionEvent;
import
org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
@@ -182,6 +183,16 @@ public class IoTDBThriftSyncConnector extends
IoTDBConnector {
return;
}
+ if (((EnrichedEvent) tabletInsertionEvent).shouldParsePattern()) {
+ if (tabletInsertionEvent instanceof PipeInsertNodeTabletInsertionEvent) {
+ transfer(
+ ((PipeInsertNodeTabletInsertionEvent)
tabletInsertionEvent).parseEventWithPattern());
+ } else { // tabletInsertionEvent instanceof PipeRawTabletInsertionEvent
+ transfer(((PipeRawTabletInsertionEvent)
tabletInsertionEvent).parseEventWithPattern());
+ }
+ return;
+ }
+
final int clientIndex = nextClientIndex();
final IoTDBThriftSyncConnectorClient client = clients.get(clientIndex);
@@ -218,6 +229,13 @@ public class IoTDBThriftSyncConnector extends
IoTDBConnector {
return;
}
+ if (((EnrichedEvent) tsFileInsertionEvent).shouldParsePattern()) {
+ for (final TabletInsertionEvent event :
tsFileInsertionEvent.toTabletInsertionEvents()) {
+ transfer(event);
+ }
+ return;
+ }
+
final int clientIndex = nextClientIndex();
final IoTDBThriftSyncConnectorClient client = clients.get(clientIndex);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/EnrichedEvent.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/EnrichedEvent.java
index ef49f2bf8c9..5c2651e5847 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/EnrichedEvent.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/EnrichedEvent.java
@@ -36,14 +36,16 @@ public abstract class EnrichedEvent implements Event {
private final AtomicInteger referenceCount;
- private final PipeTaskMeta pipeTaskMeta;
+ protected final PipeTaskMeta pipeTaskMeta;
private final String pattern;
+ private final boolean isPatternParsed;
protected EnrichedEvent(PipeTaskMeta pipeTaskMeta, String pattern) {
referenceCount = new AtomicInteger(0);
this.pipeTaskMeta = pipeTaskMeta;
this.pattern = pattern;
+ isPatternParsed =
getPattern().equals(PipeExtractorConstant.EXTRACTOR_PATTERN_DEFAULT_VALUE);
}
/**
@@ -102,7 +104,7 @@ public abstract class EnrichedEvent implements Event {
*/
public abstract boolean internallyDecreaseResourceReferenceCount(String
holderMessage);
- private void reportProgress() {
+ protected void reportProgress() {
if (pipeTaskMeta != null) {
pipeTaskMeta.updateProgressIndex(getProgressIndex());
}
@@ -128,11 +130,17 @@ public abstract class EnrichedEvent implements Event {
return pattern == null ?
PipeExtractorConstant.EXTRACTOR_PATTERN_DEFAULT_VALUE : pattern;
}
+ public boolean shouldParsePattern() {
+ return !isPatternParsed;
+ }
+
public abstract EnrichedEvent
shallowCopySelfAndBindPipeTaskMetaForProgressReport(
PipeTaskMeta pipeTaskMeta, String pattern);
public void reportException(PipeRuntimeException pipeRuntimeException) {
- PipeAgent.runtime().report(this.pipeTaskMeta, pipeRuntimeException);
+ if (pipeTaskMeta != null) {
+ PipeAgent.runtime().report(pipeTaskMeta, pipeRuntimeException);
+ }
}
public abstract boolean isGeneratedByPipe();
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/row/PipeRowCollector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/row/PipeRowCollector.java
index f67744f8b06..76c4016936f 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/row/PipeRowCollector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/row/PipeRowCollector.java
@@ -20,6 +20,8 @@
package org.apache.iotdb.db.pipe.event.common.row;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
+import org.apache.iotdb.commons.pipe.task.meta.PipeTaskMeta;
+import org.apache.iotdb.db.pipe.event.EnrichedEvent;
import
org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
import org.apache.iotdb.pipe.api.access.Row;
import org.apache.iotdb.pipe.api.collector.RowCollector;
@@ -37,6 +39,13 @@ public class PipeRowCollector implements RowCollector {
private final List<TabletInsertionEvent> tabletInsertionEventList = new
ArrayList<>();
private Tablet tablet = null;
private boolean isAligned = false;
+ private final PipeTaskMeta pipeTaskMeta; // used to report progress
+ private final EnrichedEvent sourceEvent; // used to report progress
+
+ public PipeRowCollector(PipeTaskMeta pipeTaskMeta, EnrichedEvent
sourceEvent) {
+ this.pipeTaskMeta = pipeTaskMeta;
+ this.sourceEvent = sourceEvent;
+ }
@Override
public void collectRow(Row row) {
@@ -85,13 +94,20 @@ public class PipeRowCollector implements RowCollector {
private void collectTabletInsertionEvent() {
if (tablet != null) {
- tabletInsertionEventList.add(new PipeRawTabletInsertionEvent(tablet,
isAligned));
+ tabletInsertionEventList.add(
+ new PipeRawTabletInsertionEvent(tablet, isAligned, pipeTaskMeta,
sourceEvent, false));
}
this.tablet = null;
}
public Iterable<TabletInsertionEvent> convertToTabletInsertionEvents() {
collectTabletInsertionEvent();
+
+ final int eventListSize = tabletInsertionEventList.size();
+ if (eventListSize > 0) { // The last event should report progress
+ ((PipeRawTabletInsertionEvent)
tabletInsertionEventList.get(eventListSize - 1))
+ .markAsNeedToReport();
+ }
return tabletInsertionEventList;
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java
index 739a2d88292..0d3ca5a5e98 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeInsertNodeTabletInsertionEvent.java
@@ -131,7 +131,8 @@ public class PipeInsertNodeTabletInsertionEvent extends
EnrichedEvent
public Iterable<TabletInsertionEvent> processRowByRow(BiConsumer<Row,
RowCollector> consumer) {
try {
if (dataContainer == null) {
- dataContainer = new TabletInsertionDataContainer(getInsertNode(),
getPattern());
+ dataContainer =
+ new TabletInsertionDataContainer(pipeTaskMeta, this,
getInsertNode(), getPattern());
}
return dataContainer.processRowByRow(consumer);
} catch (Exception e) {
@@ -143,7 +144,8 @@ public class PipeInsertNodeTabletInsertionEvent extends
EnrichedEvent
public Iterable<TabletInsertionEvent> processTablet(BiConsumer<Tablet,
RowCollector> consumer) {
try {
if (dataContainer == null) {
- dataContainer = new TabletInsertionDataContainer(getInsertNode(),
getPattern());
+ dataContainer =
+ new TabletInsertionDataContainer(pipeTaskMeta, this,
getInsertNode(), getPattern());
}
return dataContainer.processTablet(consumer);
} catch (Exception e) {
@@ -160,7 +162,8 @@ public class PipeInsertNodeTabletInsertionEvent extends
EnrichedEvent
public Tablet convertToTablet() {
try {
if (dataContainer == null) {
- dataContainer = new TabletInsertionDataContainer(getInsertNode(),
getPattern());
+ dataContainer =
+ new TabletInsertionDataContainer(pipeTaskMeta, this,
getInsertNode(), getPattern());
}
return dataContainer.convertToTablet();
} catch (Exception e) {
@@ -168,6 +171,12 @@ public class PipeInsertNodeTabletInsertionEvent extends
EnrichedEvent
}
}
+ /////////////////////////// parsePattern ///////////////////////////
+
+ public TabletInsertionEvent parseEventWithPattern() {
+ return new PipeRawTabletInsertionEvent(convertToTablet(), isAligned,
pipeTaskMeta, this, true);
+ }
+
/////////////////////////// Object ///////////////////////////
@Override
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
index b769be62237..5f30d3f934e 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
@@ -19,7 +19,11 @@
package org.apache.iotdb.db.pipe.event.common.tablet;
-import org.apache.iotdb.db.pipe.config.constant.PipeExtractorConstant;
+import org.apache.iotdb.commons.consensus.index.ProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex;
+import org.apache.iotdb.commons.pipe.task.meta.PipeTaskMeta;
+import org.apache.iotdb.commons.utils.TestOnly;
+import org.apache.iotdb.db.pipe.event.EnrichedEvent;
import org.apache.iotdb.pipe.api.access.Row;
import org.apache.iotdb.pipe.api.collector.RowCollector;
import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
@@ -28,26 +32,85 @@ import org.apache.iotdb.tsfile.write.record.Tablet;
import java.util.Objects;
import java.util.function.BiConsumer;
-public class PipeRawTabletInsertionEvent implements TabletInsertionEvent {
+public class PipeRawTabletInsertionEvent extends EnrichedEvent implements
TabletInsertionEvent {
private final Tablet tablet;
private final boolean isAligned;
- private final String pattern;
+
+ private final EnrichedEvent sourceEvent;
+ private boolean needToReport;
private TabletInsertionDataContainer dataContainer;
+ private PipeRawTabletInsertionEvent(
+ Tablet tablet,
+ boolean isAligned,
+ EnrichedEvent sourceEvent,
+ boolean needToReport,
+ PipeTaskMeta pipeTaskMeta,
+ String pattern) {
+ super(pipeTaskMeta, pattern);
+ this.tablet = Objects.requireNonNull(tablet);
+ this.isAligned = isAligned;
+ this.sourceEvent = sourceEvent;
+ this.needToReport = needToReport;
+ }
+
+ public PipeRawTabletInsertionEvent(
+ Tablet tablet,
+ boolean isAligned,
+ PipeTaskMeta pipeTaskMeta,
+ EnrichedEvent sourceEvent,
+ boolean needToReport) {
+ this(tablet, isAligned, sourceEvent, needToReport, pipeTaskMeta, null);
+ }
+
+ @TestOnly
public PipeRawTabletInsertionEvent(Tablet tablet, boolean isAligned) {
- this(tablet, isAligned, null);
+ this(tablet, isAligned, null, false, null, null);
}
+ @TestOnly
public PipeRawTabletInsertionEvent(Tablet tablet, boolean isAligned, String
pattern) {
- this.tablet = Objects.requireNonNull(tablet);
- this.isAligned = isAligned;
- this.pattern = pattern;
+ this(tablet, isAligned, null, false, null, pattern);
+ }
+
+ @Override
+ public boolean internallyIncreaseResourceReferenceCount(String
holderMessage) {
+ return true;
+ }
+
+ @Override
+ public boolean internallyDecreaseResourceReferenceCount(String
holderMessage) {
+ return true;
+ }
+
+ @Override
+ protected void reportProgress() {
+ if (needToReport) {
+ super.reportProgress();
+ }
+ }
+
+ @Override
+ public ProgressIndex getProgressIndex() {
+ return sourceEvent != null ? sourceEvent.getProgressIndex() : new
MinimumProgressIndex();
+ }
+
+ @Override
+ public EnrichedEvent shallowCopySelfAndBindPipeTaskMetaForProgressReport(
+ PipeTaskMeta pipeTaskMeta, String pattern) {
+ return new PipeRawTabletInsertionEvent(
+ tablet, isAligned, sourceEvent, needToReport, pipeTaskMeta, pattern);
+ }
+
+ @Override
+ public boolean isGeneratedByPipe() {
+ throw new UnsupportedOperationException("isGeneratedByPipe() is not
supported!");
}
- public String getPattern() {
- return pattern == null ?
PipeExtractorConstant.EXTRACTOR_PATTERN_DEFAULT_VALUE : pattern;
+ public void markAsNeedToReport() {
+ this.needToReport = true;
}
/////////////////////////// TabletInsertionEvent ///////////////////////////
@@ -55,7 +118,8 @@ public class PipeRawTabletInsertionEvent implements
TabletInsertionEvent {
@Override
public Iterable<TabletInsertionEvent> processRowByRow(BiConsumer<Row,
RowCollector> consumer) {
if (dataContainer == null) {
- dataContainer = new TabletInsertionDataContainer(tablet, isAligned,
getPattern());
+ dataContainer =
+ new TabletInsertionDataContainer(pipeTaskMeta, this, tablet,
isAligned, getPattern());
}
return dataContainer.processRowByRow(consumer);
}
@@ -63,7 +127,8 @@ public class PipeRawTabletInsertionEvent implements
TabletInsertionEvent {
@Override
public Iterable<TabletInsertionEvent> processTablet(BiConsumer<Tablet,
RowCollector> consumer) {
if (dataContainer == null) {
- dataContainer = new TabletInsertionDataContainer(tablet, isAligned,
getPattern());
+ dataContainer =
+ new TabletInsertionDataContainer(pipeTaskMeta, this, tablet,
isAligned, getPattern());
}
return dataContainer.processTablet(consumer);
}
@@ -75,17 +140,22 @@ public class PipeRawTabletInsertionEvent implements
TabletInsertionEvent {
}
public Tablet convertToTablet() {
- final String notNullPattern = getPattern();
-
- // if notNullPattern is "root", we don't need to convert, just return the
original tablet
- if
(notNullPattern.equals(PipeExtractorConstant.EXTRACTOR_PATTERN_DEFAULT_VALUE)) {
+ if (!shouldParsePattern()) {
return tablet;
}
// if notNullPattern is not "root", we need to convert the tablet
if (dataContainer == null) {
- dataContainer = new TabletInsertionDataContainer(tablet, isAligned,
notNullPattern);
+ dataContainer =
+ new TabletInsertionDataContainer(pipeTaskMeta, this, tablet,
isAligned, getPattern());
}
return dataContainer.convertToTablet();
}
+
+ /////////////////////////// parsePattern ///////////////////////////
+
+ public TabletInsertionEvent parseEventWithPattern() {
+ return new PipeRawTabletInsertionEvent(
+ convertToTablet(), isAligned, pipeTaskMeta, this, needToReport);
+ }
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/TabletInsertionDataContainer.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/TabletInsertionDataContainer.java
index 7fd75ac1de7..3e7506adef3 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/TabletInsertionDataContainer.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/TabletInsertionDataContainer.java
@@ -19,6 +19,8 @@
package org.apache.iotdb.db.pipe.event.common.tablet;
+import org.apache.iotdb.commons.pipe.task.meta.PipeTaskMeta;
+import org.apache.iotdb.db.pipe.event.EnrichedEvent;
import org.apache.iotdb.db.pipe.event.common.row.PipeRow;
import org.apache.iotdb.db.pipe.event.common.row.PipeRowCollector;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertNode;
@@ -45,6 +47,9 @@ import java.util.stream.IntStream;
public class TabletInsertionDataContainer {
+ private final PipeTaskMeta pipeTaskMeta; // used to report progress
+ private final EnrichedEvent sourceEvent; // used to report progress
+
private String deviceId;
private boolean isAligned;
private MeasurementSchema[] measurementSchemaList;
@@ -60,6 +65,14 @@ public class TabletInsertionDataContainer {
private Tablet tablet;
public TabletInsertionDataContainer(InsertNode insertNode, String pattern) {
+ this(null, null, insertNode, pattern);
+ }
+
+ public TabletInsertionDataContainer(
+ PipeTaskMeta pipeTaskMeta, EnrichedEvent sourceEvent, InsertNode
insertNode, String pattern) {
+ this.pipeTaskMeta = pipeTaskMeta;
+ this.sourceEvent = sourceEvent;
+
if (insertNode instanceof InsertRowNode) {
parse((InsertRowNode) insertNode, pattern);
} else if (insertNode instanceof InsertTabletNode) {
@@ -70,7 +83,15 @@ public class TabletInsertionDataContainer {
}
}
- public TabletInsertionDataContainer(Tablet tablet, boolean isAligned, String
pattern) {
+ public TabletInsertionDataContainer(
+ PipeTaskMeta pipeTaskMeta,
+ EnrichedEvent sourceEvent,
+ Tablet tablet,
+ boolean isAligned,
+ String pattern) {
+ this.pipeTaskMeta = pipeTaskMeta;
+ this.sourceEvent = sourceEvent;
+
parse(tablet, isAligned, pattern);
}
@@ -305,7 +326,7 @@ public class TabletInsertionDataContainer {
return Collections.emptyList();
}
- final PipeRowCollector rowCollector = new PipeRowCollector();
+ final PipeRowCollector rowCollector = new PipeRowCollector(pipeTaskMeta,
sourceEvent);
for (int i = 0; i < rowCount; i++) {
consumer.accept(
new PipeRow(
@@ -324,7 +345,7 @@ public class TabletInsertionDataContainer {
}
public Iterable<TabletInsertionEvent> processTablet(BiConsumer<Tablet,
RowCollector> consumer) {
- final PipeRowCollector rowCollector = new PipeRowCollector();
+ final PipeRowCollector rowCollector = new PipeRowCollector(pipeTaskMeta,
sourceEvent);
consumer.accept(convertToTablet(), rowCollector);
return rowCollector.convertToTabletInsertionEvents();
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
index 5705f5c8600..051e6d0437c 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
@@ -174,7 +174,9 @@ public class PipeTsFileInsertionEvent extends EnrichedEvent
implements TsFileIns
try {
if (dataContainer == null) {
waitForTsFileClose();
- dataContainer = new TsFileInsertionDataContainer(tsFile, getPattern(),
startTime, endTime);
+ dataContainer =
+ new TsFileInsertionDataContainer(
+ tsFile, getPattern(), startTime, endTime, pipeTaskMeta, this);
}
return dataContainer.toTabletInsertionEvents();
} catch (InterruptedException e) {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/TsFileInsertionDataContainer.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/TsFileInsertionDataContainer.java
index 0d72fd06e53..6c118880e77 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/TsFileInsertionDataContainer.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/TsFileInsertionDataContainer.java
@@ -19,6 +19,8 @@
package org.apache.iotdb.db.pipe.event.common.tsfile;
+import org.apache.iotdb.commons.pipe.task.meta.PipeTaskMeta;
+import org.apache.iotdb.db.pipe.event.EnrichedEvent;
import
org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
import org.apache.iotdb.pipe.api.exception.PipeException;
@@ -50,9 +52,11 @@ public class TsFileInsertionDataContainer implements
AutoCloseable {
private static final Logger LOGGER =
LoggerFactory.getLogger(TsFileInsertionDataContainer.class);
- // used to filter data
- private final String pattern;
- private final IExpression timeFilterExpression;
+ private final String pattern; // used to filter data
+ private final IExpression timeFilterExpression; // used to filter data
+
+ private final PipeTaskMeta pipeTaskMeta; // used to report progress
+ private final EnrichedEvent sourceEvent; // used to report progress
private final TsFileSequenceReader tsFileSequenceReader;
private final TsFileReader tsFileReader;
@@ -63,6 +67,17 @@ public class TsFileInsertionDataContainer implements
AutoCloseable {
public TsFileInsertionDataContainer(File tsFile, String pattern, long
startTime, long endTime)
throws IOException {
+ this(tsFile, pattern, startTime, endTime, null, null);
+ }
+
+ public TsFileInsertionDataContainer(
+ File tsFile,
+ String pattern,
+ long startTime,
+ long endTime,
+ PipeTaskMeta pipeTaskMeta,
+ EnrichedEvent sourceEvent)
+ throws IOException {
this.pattern = pattern;
timeFilterExpression =
(startTime == Long.MIN_VALUE && endTime == Long.MAX_VALUE)
@@ -71,6 +86,9 @@ public class TsFileInsertionDataContainer implements
AutoCloseable {
new GlobalTimeExpression(TimeFilter.gtEq(startTime)),
new GlobalTimeExpression(TimeFilter.ltEq(endTime)));
+ this.pipeTaskMeta = pipeTaskMeta;
+ this.sourceEvent = sourceEvent;
+
try {
tsFileSequenceReader = new
TsFileSequenceReader(tsFile.getAbsolutePath());
tsFileReader = new TsFileReader(tsFileSequenceReader);
@@ -176,12 +194,18 @@ public class TsFileInsertionDataContainer implements
AutoCloseable {
final Tablet tablet = tabletIterator.next();
final boolean isAligned =
deviceIsAlignedMap.getOrDefault(tablet.deviceId, false);
- final TabletInsertionEvent next = new
PipeRawTabletInsertionEvent(tablet, isAligned);
+ final TabletInsertionEvent next;
if (!hasNext()) {
+ next =
+ new PipeRawTabletInsertionEvent(
+ tablet, isAligned, pipeTaskMeta, sourceEvent, true);
close();
+ } else {
+ next =
+ new PipeRawTabletInsertionEvent(
+ tablet, isAligned, pipeTaskMeta, sourceEvent, false);
}
-
return next;
}
};
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/PipeDoNothingProcessor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/PipeDoNothingProcessor.java
index 7de66c55ca9..e1de3e13331 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/PipeDoNothingProcessor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/PipeDoNothingProcessor.java
@@ -19,9 +19,6 @@
package org.apache.iotdb.db.pipe.processor;
-import org.apache.iotdb.db.pipe.config.constant.PipeExtractorConstant;
-import org.apache.iotdb.db.pipe.event.EnrichedEvent;
-import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
import org.apache.iotdb.pipe.api.PipeProcessor;
import org.apache.iotdb.pipe.api.collector.EventCollector;
import
org.apache.iotdb.pipe.api.customizer.configuration.PipeProcessorRuntimeConfiguration;
@@ -30,7 +27,6 @@ import
org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
import org.apache.iotdb.pipe.api.event.Event;
import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
import org.apache.iotdb.pipe.api.event.dml.insertion.TsFileInsertionEvent;
-import org.apache.iotdb.pipe.api.exception.PipeException;
import java.io.IOException;
@@ -50,54 +46,13 @@ public class PipeDoNothingProcessor implements
PipeProcessor {
@Override
public void process(TabletInsertionEvent tabletInsertionEvent,
EventCollector eventCollector)
throws IOException {
- if (tabletInsertionEvent instanceof EnrichedEvent) {
- final EnrichedEvent enrichedEvent = (EnrichedEvent) tabletInsertionEvent;
- if (enrichedEvent
- .getPattern()
- .equals(PipeExtractorConstant.EXTRACTOR_PATTERN_DEFAULT_VALUE)) {
- eventCollector.collect(tabletInsertionEvent);
- } else {
- tabletInsertionEvent
- .processRowByRow(
- (row, rowCollector) -> {
- try {
- rowCollector.collectRow(row);
- } catch (IOException e) {
- throw new PipeException("Failed to collect row", e);
- }
- })
- .forEach(
- event -> {
- try {
- eventCollector.collect(event);
- } catch (IOException e) {
- throw new PipeException("Failed to collect event", e);
- }
- });
- }
- } else {
- eventCollector.collect(tabletInsertionEvent);
- }
+ eventCollector.collect(tabletInsertionEvent);
}
@Override
public void process(TsFileInsertionEvent tsFileInsertionEvent,
EventCollector eventCollector)
throws IOException {
- if (tsFileInsertionEvent instanceof PipeTsFileInsertionEvent) {
- final PipeTsFileInsertionEvent enrichedEvent =
- (PipeTsFileInsertionEvent) tsFileInsertionEvent;
- if
(enrichedEvent.getPattern().equals(PipeExtractorConstant.EXTRACTOR_PATTERN_DEFAULT_VALUE)
- && !enrichedEvent.hasTimeFilter()) {
- eventCollector.collect(tsFileInsertionEvent);
- } else {
- for (final TabletInsertionEvent tabletInsertionEvent :
- tsFileInsertionEvent.toTabletInsertionEvents()) {
- eventCollector.collect(tabletInsertionEvent);
- }
- }
- } else {
- eventCollector.collect(tsFileInsertionEvent);
- }
+ eventCollector.collect(tsFileInsertionEvent);
}
@Override
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/PipeEventCollector.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/PipeEventCollector.java
index 19ef3f6fd2f..8ec84529d08 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/PipeEventCollector.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/PipeEventCollector.java
@@ -20,7 +20,6 @@
package org.apache.iotdb.db.pipe.task.connection;
import org.apache.iotdb.db.pipe.event.EnrichedEvent;
-import
org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
import org.apache.iotdb.pipe.api.collector.EventCollector;
import org.apache.iotdb.pipe.api.event.Event;
@@ -56,37 +55,13 @@ public class PipeEventCollector implements EventCollector {
if (pendingQueue.waitedOffer(bufferedEvent)) {
bufferQueue.poll();
} else {
- // If timeout, we judge whether the new event is a
PipeRawTabletInsertionEvent. If it is,
- // we wait for pending queue to be available without timeout until the
pending queue is
- // available. We don't put PipeRawTabletInsertionEvent into buffer
queue, because it is
- // memory consuming, holding too many PipeRawTabletInsertionEvent in
buffer queue may cause
- // OOM.
- if (event instanceof PipeRawTabletInsertionEvent) {
- if (pendingQueue.put(bufferedEvent)) {
- bufferQueue.poll();
- } else {
- LOGGER.warn("interrupted when putting event into pending queue,
event: {}", event);
- bufferQueue.offer(event);
- return;
- }
- } else {
- bufferQueue.offer(event);
- return;
- }
+ bufferQueue.offer(event);
+ return;
}
}
if (!pendingQueue.waitedOffer(event)) {
- // PipeRawTabletInsertionEvent is memory consuming, so we should not put
it into buffer queue
- // when pending queue is full. Otherwise, it may cause OOM.
- if (event instanceof PipeRawTabletInsertionEvent) {
- if (!pendingQueue.put(event)) {
- LOGGER.warn("interrupted when putting event into pending queue,
event: {}", event);
- bufferQueue.offer(event);
- }
- } else {
- bufferQueue.offer(event);
- }
+ bufferQueue.offer(event);
}
}
}