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);
     }
   }
 }


Reply via email to