This is an automated email from the ASF dual-hosted git repository.

Caideyipi pushed a commit to branch fix/builtin-sink-synchronous-drop
in repository https://gitbox.apache.org/repos/asf/iotdb.git

commit 420ecd71861e1e07f9f7e6bebb730ee93b7ce0be
Author: Caideyipi <[email protected]>
AuthorDate: Wed Sep 16 11:06:21 2026 +0800

    [Pipe] Drop builtin sinks synchronously
---
 .../agent/task/subtask/sink/PipeSinkSubtask.java   |  59 ++++++++--
 .../task/subtask/sink/PipeSinkSubtaskManager.java  |   4 +-
 .../task/subtask/SubscriptionSinkSubtask.java      |   5 +-
 .../task/subtask/sink/PipeSinkSubtaskTest.java     | 123 ++++++++++++++++++++-
 .../agent/plugin/builtin/BuiltinPipePlugin.java    |  33 ++++++
 5 files changed, 211 insertions(+), 13 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtask.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtask.java
index b0de861e41b..9cd0eb59f38 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtask.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtask.java
@@ -72,6 +72,7 @@ public class PipeSinkSubtask extends PipeAbstractSinkSubtask {
   private final String attributeSortedString;
   private final String attributeDisplayString;
   private final int sinkIndex;
+  private final boolean isExternalSink;
 
   // Now parallel connectors run the same time, thus the heartbeat events are 
not sure
   // to trigger the general event transfer function, causing potentially such 
as
@@ -96,7 +97,8 @@ public class PipeSinkSubtask extends PipeAbstractSinkSubtask {
         attributeSortedString,
         sinkIndex,
         inputPendingQueue,
-        outputPipeConnector);
+        outputPipeConnector,
+        true);
   }
 
   public PipeSinkSubtask(
@@ -115,7 +117,8 @@ public class PipeSinkSubtask extends 
PipeAbstractSinkSubtask {
         attributeSortedString,
         sinkIndex,
         inputPendingQueue,
-        outputPipeConnector);
+        outputPipeConnector,
+        true);
   }
 
   public PipeSinkSubtask(
@@ -134,7 +137,8 @@ public class PipeSinkSubtask extends 
PipeAbstractSinkSubtask {
         attributeDisplayString,
         sinkIndex,
         inputPendingQueue,
-        outputPipeConnector);
+        outputPipeConnector,
+        true);
   }
 
   public PipeSinkSubtask(
@@ -146,11 +150,34 @@ public class PipeSinkSubtask extends 
PipeAbstractSinkSubtask {
       final int sinkIndex,
       final UnboundedBlockingPendingQueue<Event> inputPendingQueue,
       final PipeConnector outputPipeConnector) {
+    this(
+        pipeName,
+        taskID,
+        creationTime,
+        attributeSortedString,
+        attributeDisplayString,
+        sinkIndex,
+        inputPendingQueue,
+        outputPipeConnector,
+        true);
+  }
+
+  public PipeSinkSubtask(
+      final String pipeName,
+      final String taskID,
+      final long creationTime,
+      final String attributeSortedString,
+      final String attributeDisplayString,
+      final int sinkIndex,
+      final UnboundedBlockingPendingQueue<Event> inputPendingQueue,
+      final PipeConnector outputPipeConnector,
+      final boolean isExternalSink) {
     super(taskID, creationTime, outputPipeConnector);
     this.pipeName = pipeName;
     this.attributeSortedString = attributeSortedString;
     this.attributeDisplayString = attributeDisplayString;
     this.sinkIndex = sinkIndex;
+    this.isExternalSink = isExternalSink;
     this.inputPendingQueue = inputPendingQueue;
 
     if (!attributeSortedString.startsWith("schema_")) {
@@ -384,6 +411,17 @@ public class PipeSinkSubtask extends 
PipeAbstractSinkSubtask {
   }
 
   private boolean closeOutputPipeSink() throws Exception {
+    if (!isExternalSink) {
+      outputPipeSinkOperationLock.lock();
+      try {
+        discardPendingEventsOfPipeUnderLock();
+        outputPipeSink.close();
+      } finally {
+        outputPipeSinkOperationLock.unlock();
+      }
+      return true;
+    }
+
     final AtomicReference<Exception> exception = new AtomicReference<>();
     final AtomicBoolean closeStarted = new AtomicBoolean(false);
     final Thread closeThread =
@@ -493,12 +531,17 @@ public class PipeSinkSubtask extends 
PipeAbstractSinkSubtask {
     }
 
     pendingDiscardCommitterKeys.offer(committerKey);
-    if (outputPipeSinkOperationLock.tryLock()) {
-      try {
-        discardPendingEventsOfPipeUnderLock();
-      } finally {
-        outputPipeSinkOperationLock.unlock();
+    if (isExternalSink) {
+      if (!outputPipeSinkOperationLock.tryLock()) {
+        return;
       }
+    } else {
+      outputPipeSinkOperationLock.lock();
+    }
+    try {
+      discardPendingEventsOfPipeUnderLock();
+    } finally {
+      outputPipeSinkOperationLock.unlock();
     }
   }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskManager.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskManager.java
index c6bc5d0c28f..7ac1f7ba12a 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskManager.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskManager.java
@@ -20,6 +20,7 @@
 package org.apache.iotdb.db.pipe.agent.task.subtask.sink;
 
 import org.apache.iotdb.commons.consensus.DataRegionId;
+import org.apache.iotdb.commons.pipe.agent.plugin.builtin.BuiltinPipePlugin;
 import 
org.apache.iotdb.commons.pipe.agent.task.connection.UnboundedBlockingPendingQueue;
 import org.apache.iotdb.commons.pipe.agent.task.meta.PipeRuntimeMeta;
 import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey;
@@ -164,7 +165,8 @@ public class PipeSinkSubtaskManager {
                 attributeDisplayStringWithPrefix,
                 sinkIndex,
                 pendingQueue,
-                pipeSink);
+                pipeSink,
+                !BuiltinPipePlugin.BUILTIN_SINKS.contains(connectorKey));
         final PipeSinkSubtaskLifeCycle pipeSinkSubtaskLifeCycle =
             new PipeSinkSubtaskLifeCycle(executor, pipeSinkSubtask, 
pendingQueue);
         pipeSinkSubtaskLifeCycleList.add(pipeSinkSubtaskLifeCycle);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionSinkSubtask.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionSinkSubtask.java
index 6c26e0fab10..586bb39d579 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionSinkSubtask.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/task/subtask/SubscriptionSinkSubtask.java
@@ -49,12 +49,15 @@ public class SubscriptionSinkSubtask extends 
PipeSinkSubtask {
       final String topicName,
       final String consumerGroupId) {
     super(
+        null,
         taskID,
         creationTime,
         attributeSortedString,
+        attributeSortedString,
         connectorIndex,
         inputPendingQueue,
-        outputPipeConnector);
+        outputPipeConnector,
+        false);
     this.topicName = topicName;
     this.consumerGroupId = consumerGroupId;
   }
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskTest.java
index 3d5427d39b3..eac860f8763 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeSinkSubtaskTest.java
@@ -89,7 +89,7 @@ public class PipeSinkSubtaskTest {
   }
 
   @Test
-  public void testDiscardEventsOfPipeNotBlockedByConnectionRetry() throws 
Exception {
+  public void testExternalSinkDiscardEventsOfPipeNotBlockedByConnectionRetry() 
throws Exception {
     final CountDownLatch handshakeEntered = new CountDownLatch(1);
     final CountDownLatch releaseHandshake = new CountDownLatch(1);
     final CountDownLatch discardEntered = new CountDownLatch(1);
@@ -145,7 +145,68 @@ public class PipeSinkSubtaskTest {
   }
 
   @Test
-  public void testCloseNotConcurrentWithConnectionRetry() throws Exception {
+  public void testBuiltinSinkDiscardEventsOfPipeWaitsForConnectionRetry() 
throws Exception {
+    final CountDownLatch handshakeEntered = new CountDownLatch(1);
+    final CountDownLatch releaseHandshake = new CountDownLatch(1);
+    final CountDownLatch discardEntered = new CountDownLatch(1);
+    final AtomicBoolean discardDuringHandshake = new AtomicBoolean(false);
+    final PipeConnector connector =
+        new BlockingHandshakeConnector(
+            handshakeEntered,
+            releaseHandshake,
+            new CountDownLatch(0),
+            new AtomicBoolean(false),
+            discardEntered,
+            discardDuringHandshake);
+    final UnboundedBlockingPendingQueue<?> pendingQueue = 
mock(UnboundedBlockingPendingQueue.class);
+
+    final PipeSinkSubtask subtask =
+        new PipeSinkSubtask(
+            null,
+            "PipeSinkSubtaskTest",
+            System.currentTimeMillis(),
+            "data_test",
+            "data_test",
+            0,
+            (UnboundedBlockingPendingQueue) pendingQueue,
+            connector,
+            false);
+
+    final Thread failureThread =
+        new Thread(() -> subtask.onFailure(new 
PipeConnectionException("connection broken")));
+    failureThread.start();
+    Assert.assertTrue(handshakeEntered.await(5, TimeUnit.SECONDS));
+
+    final CountDownLatch discardReturned = new CountDownLatch(1);
+    final Thread discardThread =
+        new Thread(
+            () -> {
+              try {
+                subtask.discardEventsOfPipe(new CommitterKey("pipe", 1L, 1, 
-1));
+              } finally {
+                discardReturned.countDown();
+              }
+            });
+    discardThread.start();
+
+    try {
+      Assert.assertFalse(discardReturned.await(100, TimeUnit.MILLISECONDS));
+      Assert.assertEquals(1L, discardEntered.getCount());
+      Assert.assertFalse(discardDuringHandshake.get());
+    } finally {
+      releaseHandshake.countDown();
+      discardThread.join(5000);
+      failureThread.join(5000);
+      subtask.close();
+    }
+
+    Assert.assertFalse(discardThread.isAlive());
+    Assert.assertTrue(discardEntered.await(1, TimeUnit.SECONDS));
+    Assert.assertFalse(discardDuringHandshake.get());
+  }
+
+  @Test
+  public void testExternalSinkCloseNotConcurrentWithConnectionRetry() throws 
Exception {
     final int originalTimeout =
         
CommonDescriptor.getInstance().getConfig().getDnConnectionTimeoutInMS();
     CommonDescriptor.getInstance().getConfig().setDnConnectionTimeoutInMS(30);
@@ -197,7 +258,7 @@ public class PipeSinkSubtaskTest {
   }
 
   @Test
-  public void testCloseDoesNotWaitForeverForConnectorClose() throws Exception {
+  public void testExternalSinkCloseDoesNotWaitForeverForConnectorClose() 
throws Exception {
     final int originalTimeout =
         
CommonDescriptor.getInstance().getConfig().getDnConnectionTimeoutInMS();
     CommonDescriptor.getInstance().getConfig().setDnConnectionTimeoutInMS(30);
@@ -238,6 +299,62 @@ public class PipeSinkSubtaskTest {
     }
   }
 
+  @Test
+  public void testBuiltinSinkCloseWaitsForConnectorClose() throws Exception {
+    final int originalTimeout =
+        
CommonDescriptor.getInstance().getConfig().getDnConnectionTimeoutInMS();
+    CommonDescriptor.getInstance().getConfig().setDnConnectionTimeoutInMS(30);
+
+    final PipeConnector connector = mock(PipeConnector.class);
+    final UnboundedBlockingPendingQueue<?> pendingQueue = 
mock(UnboundedBlockingPendingQueue.class);
+    final CountDownLatch closeEntered = new CountDownLatch(1);
+    final CountDownLatch releaseClose = new CountDownLatch(1);
+
+    doAnswer(
+            invocation -> {
+              closeEntered.countDown();
+              releaseClose.await(5, TimeUnit.SECONDS);
+              return null;
+            })
+        .when(connector)
+        .close();
+
+    final PipeSinkSubtask subtask =
+        new PipeSinkSubtask(
+            null,
+            "PipeSinkSubtaskTest",
+            System.currentTimeMillis(),
+            "data_test",
+            "data_test",
+            0,
+            (UnboundedBlockingPendingQueue) pendingQueue,
+            connector,
+            false);
+    final CountDownLatch closeReturned = new CountDownLatch(1);
+    final Thread closeThread =
+        new Thread(
+            () -> {
+              try {
+                subtask.close();
+              } finally {
+                closeReturned.countDown();
+              }
+            });
+
+    try {
+      closeThread.start();
+      Assert.assertTrue(closeEntered.await(5, TimeUnit.SECONDS));
+      Assert.assertFalse(closeReturned.await(100, TimeUnit.MILLISECONDS));
+    } finally {
+      releaseClose.countDown();
+      closeThread.join(5000);
+      
CommonDescriptor.getInstance().getConfig().setDnConnectionTimeoutInMS(originalTimeout);
+    }
+
+    Assert.assertFalse(closeThread.isAlive());
+    Assert.assertEquals(0L, closeReturned.getCount());
+  }
+
   @Test
   public void testTransferExceptionUsesDisplayTaskID() throws Exception {
     final PipeConnector connector = mock(PipeConnector.class);
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/plugin/builtin/BuiltinPipePlugin.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/plugin/builtin/BuiltinPipePlugin.java
index 76b45d1e798..7fd99dbc2aa 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/plugin/builtin/BuiltinPipePlugin.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/plugin/builtin/BuiltinPipePlugin.java
@@ -145,6 +145,39 @@ public enum BuiltinPipePlugin {
                   DO_NOTHING_SOURCE.getPipePluginName().toLowerCase(),
                   IOTDB_SOURCE.getPipePluginName().toLowerCase())));
 
+  // Used to distinguish between builtin and external sinks.
+  public static final Set<String> BUILTIN_SINKS =
+      Collections.unmodifiableSet(
+          new HashSet<>(
+              Arrays.asList(
+                  DO_NOTHING_CONNECTOR.getPipePluginName().toLowerCase(),
+                  IOTDB_THRIFT_CONNECTOR.getPipePluginName().toLowerCase(),
+                  IOTDB_THRIFT_SSL_CONNECTOR.getPipePluginName().toLowerCase(),
+                  
IOTDB_THRIFT_SYNC_CONNECTOR.getPipePluginName().toLowerCase(),
+                  
IOTDB_THRIFT_ASYNC_CONNECTOR.getPipePluginName().toLowerCase(),
+                  
IOTDB_LEGACY_PIPE_CONNECTOR.getPipePluginName().toLowerCase(),
+                  IOTDB_AIR_GAP_CONNECTOR.getPipePluginName().toLowerCase(),
+                  
IOT_CONSENSUS_V2_ASYNC_CONNECTOR.getPipePluginName().toLowerCase(),
+                  
PIPE_CONSENSUS_ASYNC_CONNECTOR.getPipePluginName().toLowerCase(),
+                  WEBSOCKET_CONNECTOR.getPipePluginName().toLowerCase(),
+                  OPC_UA_CONNECTOR.getPipePluginName().toLowerCase(),
+                  OPC_DA_CONNECTOR.getPipePluginName().toLowerCase(),
+                  WRITE_BACK_CONNECTOR.getPipePluginName().toLowerCase(),
+                  DO_NOTHING_SINK.getPipePluginName().toLowerCase(),
+                  IOTDB_THRIFT_SINK.getPipePluginName().toLowerCase(),
+                  IOTDB_THRIFT_SSL_SINK.getPipePluginName().toLowerCase(),
+                  IOTDB_THRIFT_SYNC_SINK.getPipePluginName().toLowerCase(),
+                  IOTDB_THRIFT_ASYNC_SINK.getPipePluginName().toLowerCase(),
+                  IOTDB_LEGACY_PIPE_SINK.getPipePluginName().toLowerCase(),
+                  IOTDB_AIR_GAP_SINK.getPipePluginName().toLowerCase(),
+                  WEBSOCKET_SINK.getPipePluginName().toLowerCase(),
+                  OPC_UA_SINK.getPipePluginName().toLowerCase(),
+                  OPC_DA_SINK.getPipePluginName().toLowerCase(),
+                  WRITE_BACK_SINK.getPipePluginName().toLowerCase(),
+                  SUBSCRIPTION_SINK.getPipePluginName().toLowerCase(),
+                  
IOT_CONSENSUS_V2_ASYNC_SINK.getPipePluginName().toLowerCase(),
+                  
PIPE_CONSENSUS_ASYNC_SINK.getPipePluginName().toLowerCase())));
+
   public static final Set<String> SHOW_PIPE_PLUGINS_BLACKLIST =
       Collections.unmodifiableSet(
           new HashSet<>(

Reply via email to