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 f1a4012424f Pipe: drop HeartbeatEvent when pipe is stopped to avoid 
OOM (#10977)
f1a4012424f is described below

commit f1a4012424fef30509ed13e33e1b641ffc5d4ba9
Author: Steve Yurong Su <[email protected]>
AuthorDate: Mon Aug 28 23:24:25 2023 +0800

    Pipe: drop HeartbeatEvent when pipe is stopped to avoid OOM (#10977)
---
 .../realtime/PipeRealtimeDataRegionHybridExtractor.java     |  9 ++++++++-
 .../realtime/PipeRealtimeDataRegionLogExtractor.java        |  9 ++++++++-
 .../realtime/PipeRealtimeDataRegionTsFileExtractor.java     |  7 +++++++
 .../iotdb/db/pipe/task/connection/BlockingPendingQueue.java |  2 +-
 .../iotdb/db/pipe/task/connection/PipeEventCollector.java   | 13 ++++++++++---
 .../pipe/task/connection/UnboundedBlockingPendingQueue.java | 12 ++++++++++--
 6 files changed, 44 insertions(+), 8 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionHybridExtractor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionHybridExtractor.java
index b9e8f275e28..edbc14e26d4 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionHybridExtractor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionHybridExtractor.java
@@ -155,11 +155,18 @@ public class PipeRealtimeDataRegionHybridExtractor 
extends PipeRealtimeDataRegio
   }
 
   private void extractHeartbeat(PipeRealtimeEvent event) {
+    if (pendingQueue.peekLast() instanceof PipeHeartbeatEvent) {
+      // if the last event in the pending queue is a heartbeat event, we 
should not extract any more
+      // heartbeat events to avoid OOM when the pipe is stopped.
+      
event.decreaseReferenceCount(PipeRealtimeDataRegionHybridExtractor.class.getName());
+      return;
+    }
+
     if (!pendingQueue.waitedOffer(event)) {
       // this would not happen, but just in case.
       // pendingQueue is unbounded, so it should never reach capacity.
       LOGGER.error(
-          "extract: pending queue of PipeRealtimeDataRegionTsFileExtractor {} "
+          "extract: pending queue of PipeRealtimeDataRegionHybridExtractor {} "
               + "has reached capacity, discard heartbeat event {}",
           this,
           event);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionLogExtractor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionLogExtractor.java
index 14e899300b8..fb3783dc224 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionLogExtractor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionLogExtractor.java
@@ -75,11 +75,18 @@ public class PipeRealtimeDataRegionLogExtractor extends 
PipeRealtimeDataRegionEx
   }
 
   private void extractHeartbeat(PipeRealtimeEvent event) {
+    if (pendingQueue.peekLast() instanceof PipeHeartbeatEvent) {
+      // if the last event in the pending queue is a heartbeat event, we 
should not extract any more
+      // heartbeat events to avoid OOM when the pipe is stopped.
+      
event.decreaseReferenceCount(PipeRealtimeDataRegionLogExtractor.class.getName());
+      return;
+    }
+
     if (!pendingQueue.waitedOffer(event)) {
       // this would not happen, but just in case.
       // pendingQueue is unbounded, so it should never reach capacity.
       LOGGER.error(
-          "extract: pending queue of PipeRealtimeDataRegionTsFileExtractor {} "
+          "extract: pending queue of PipeRealtimeDataRegionLogExtractor {} "
               + "has reached capacity, discard heartbeat event {}",
           this,
           event);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionTsFileExtractor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionTsFileExtractor.java
index 9546e35906d..5967b1bc87b 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionTsFileExtractor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/realtime/PipeRealtimeDataRegionTsFileExtractor.java
@@ -75,6 +75,13 @@ public class PipeRealtimeDataRegionTsFileExtractor extends 
PipeRealtimeDataRegio
   }
 
   private void extractHeartbeat(PipeRealtimeEvent event) {
+    if (pendingQueue.peekLast() instanceof PipeHeartbeatEvent) {
+      // if the last event in the pending queue is a heartbeat event, we 
should not extract any more
+      // heartbeat events to avoid OOM when the pipe is stopped.
+      
event.decreaseReferenceCount(PipeRealtimeDataRegionTsFileExtractor.class.getName());
+      return;
+    }
+
     if (!pendingQueue.waitedOffer(event)) {
       // This would not happen, but just in case.
       // Pending is unbounded, so it should never reach capacity.
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/BlockingPendingQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/BlockingPendingQueue.java
index f81d301b6f2..fc7e2cd6f5c 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/BlockingPendingQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/BlockingPendingQueue.java
@@ -35,7 +35,7 @@ public abstract class BlockingPendingQueue<E extends Event> {
   private static final long MAX_BLOCKING_TIME_MS =
       
PipeConfig.getInstance().getPipeSubtaskExecutorPendingQueueMaxBlockingTimeMs();
 
-  private final BlockingQueue<E> pendingQueue;
+  protected final BlockingQueue<E> pendingQueue;
 
   protected BlockingPendingQueue(BlockingQueue<E> pendingQueue) {
     this.pendingQueue = pendingQueue;
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 bf57908c71b..772851305ee 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,17 +20,18 @@
 package org.apache.iotdb.db.pipe.task.connection;
 
 import org.apache.iotdb.db.pipe.event.EnrichedEvent;
+import org.apache.iotdb.db.pipe.event.common.heartbeat.PipeHeartbeatEvent;
 import org.apache.iotdb.pipe.api.collector.EventCollector;
 import org.apache.iotdb.pipe.api.event.Event;
 
+import java.util.Deque;
 import java.util.LinkedList;
-import java.util.Queue;
 
 public class PipeEventCollector implements EventCollector {
 
   private final BoundedBlockingPendingQueue<Event> pendingQueue;
 
-  private final Queue<Event> bufferQueue;
+  private final Deque<Event> bufferQueue;
 
   public PipeEventCollector(BoundedBlockingPendingQueue<Event> pendingQueue) {
     this.pendingQueue = pendingQueue;
@@ -50,7 +51,13 @@ public class PipeEventCollector implements EventCollector {
       if (pendingQueue.waitedOffer(bufferedEvent)) {
         bufferQueue.poll();
       } else {
-        bufferQueue.offer(event);
+        // We can NOT keep too many PipeHeartbeatEvent in bufferQueue because 
they may cause OOM.
+        if (event instanceof PipeHeartbeatEvent
+            && bufferQueue.peekLast() instanceof PipeHeartbeatEvent) {
+          ((EnrichedEvent) 
event).decreaseReferenceCount(PipeEventCollector.class.getName());
+        } else {
+          bufferQueue.offer(event);
+        }
         return;
       }
     }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/UnboundedBlockingPendingQueue.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/UnboundedBlockingPendingQueue.java
index dafb567e902..343621bbb4a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/UnboundedBlockingPendingQueue.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/task/connection/UnboundedBlockingPendingQueue.java
@@ -21,11 +21,19 @@ package org.apache.iotdb.db.pipe.task.connection;
 
 import org.apache.iotdb.pipe.api.event.Event;
 
-import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.BlockingDeque;
+import java.util.concurrent.LinkedBlockingDeque;
 
 public class UnboundedBlockingPendingQueue<E extends Event> extends 
BlockingPendingQueue<E> {
 
+  private final BlockingDeque<E> pendingDeque;
+
   public UnboundedBlockingPendingQueue() {
-    super(new LinkedBlockingQueue<>());
+    super(new LinkedBlockingDeque<>());
+    pendingDeque = (BlockingDeque<E>) pendingQueue;
+  }
+
+  public E peekLast() {
+    return pendingDeque.peekLast();
   }
 }

Reply via email to