This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/dev/1.3 by this push:
new e1f402f3e29 [Pipe] Drop builtin sinks synchronously (#18646) (#18665)
e1f402f3e29 is described below
commit e1f402f3e29d381a1aceca30853fb65baa16f7d1
Author: Caideyipi <[email protected]>
AuthorDate: Mon Sep 21 11:43:53 2026 +0800
[Pipe] Drop builtin sinks synchronously (#18646) (#18665)
(cherry picked from commit e967831ca2a7922983e6dd1b27c2d41031574b9d)
---
.../agent/task/subtask/sink/PipeSinkSubtask.java | 51 ++++++++-
.../task/subtask/sink/PipeSinkSubtaskManager.java | 4 +-
.../task/subtask/SubscriptionSinkSubtask.java | 4 +-
.../task/subtask/sink/PipeSinkSubtaskTest.java | 121 ++++++++++++++++++++-
.../agent/plugin/builtin/BuiltinPipePlugin.java | 31 ++++++
5 files changed, 200 insertions(+), 11 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 0046a9aded8..c6607b46a1e 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
@@ -69,6 +69,7 @@ public class PipeSinkSubtask extends PipeAbstractSinkSubtask {
private final String pipeName;
private final String attributeSortedString;
private final int connectorIndex;
+ 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
@@ -93,7 +94,8 @@ public class PipeSinkSubtask extends PipeAbstractSinkSubtask {
attributeSortedString,
connectorIndex,
inputPendingQueue,
- outputPipeConnector);
+ outputPipeConnector,
+ true);
}
public PipeSinkSubtask(
@@ -104,10 +106,31 @@ public class PipeSinkSubtask extends
PipeAbstractSinkSubtask {
final int connectorIndex,
final UnboundedBlockingPendingQueue<Event> inputPendingQueue,
final PipeConnector outputPipeConnector) {
+ this(
+ pipeName,
+ taskID,
+ creationTime,
+ attributeSortedString,
+ connectorIndex,
+ inputPendingQueue,
+ outputPipeConnector,
+ true);
+ }
+
+ public PipeSinkSubtask(
+ final String pipeName,
+ final String taskID,
+ final long creationTime,
+ final String attributeSortedString,
+ final int connectorIndex,
+ final UnboundedBlockingPendingQueue<Event> inputPendingQueue,
+ final PipeConnector outputPipeConnector,
+ final boolean isExternalSink) {
super(taskID, creationTime, outputPipeConnector);
this.pipeName = pipeName;
this.attributeSortedString = attributeSortedString;
this.connectorIndex = connectorIndex;
+ this.isExternalSink = isExternalSink;
this.inputPendingQueue = inputPendingQueue;
if (!attributeSortedString.startsWith("schema_")) {
@@ -291,6 +314,17 @@ public class PipeSinkSubtask extends
PipeAbstractSinkSubtask {
}
private boolean closeOutputPipeConnector() throws Exception {
+ if (!isExternalSink) {
+ outputPipeConnectorOperationLock.lock();
+ try {
+ discardPendingEventsOfPipeUnderLock();
+ outputPipeConnector.close();
+ } finally {
+ outputPipeConnectorOperationLock.unlock();
+ }
+ return true;
+ }
+
final AtomicReference<Exception> exception = new AtomicReference<>();
final AtomicBoolean closeStarted = new AtomicBoolean(false);
final Thread closeThread =
@@ -397,12 +431,17 @@ public class PipeSinkSubtask extends
PipeAbstractSinkSubtask {
}
pendingDiscardCommitterKeys.offer(committerKey);
- if (outputPipeConnectorOperationLock.tryLock()) {
- try {
- discardPendingEventsOfPipeUnderLock();
- } finally {
- outputPipeConnectorOperationLock.unlock();
+ if (isExternalSink) {
+ if (!outputPipeConnectorOperationLock.tryLock()) {
+ return;
}
+ } else {
+ outputPipeConnectorOperationLock.lock();
+ }
+ try {
+ discardPendingEventsOfPipeUnderLock();
+ } finally {
+ outputPipeConnectorOperationLock.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 dea65d86ee0..c28089090c7 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;
@@ -153,7 +154,8 @@ public class PipeSinkSubtaskManager {
attributeSortedString,
connectorIndex,
pendingQueue,
- pipeConnector);
+ pipeConnector,
+ !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 2ca332263b5..91363924c7b 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
@@ -47,12 +47,14 @@ public class SubscriptionSinkSubtask extends
PipeSinkSubtask {
final String topicName,
final String consumerGroupId) {
super(
+ null,
taskID,
creationTime,
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 6c5e79ff68b..2b99cccbe9a 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
@@ -88,7 +88,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);
@@ -143,7 +143,67 @@ 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",
+ 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);
@@ -194,7 +254,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);
@@ -234,6 +294,61 @@ 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",
+ 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
testTransferExceptionWithNullRootCauseMessageIncludesExceptionType()
throws Exception {
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 e5c041e8ec6..c18e0b7caa7 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
@@ -136,6 +136,37 @@ 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(),
+
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(),
+
PIPE_CONSENSUS_ASYNC_SINK.getPipePluginName().toLowerCase())));
+
public static final Set<String> SHOW_PIPE_PLUGINS_BLACKLIST =
Collections.unmodifiableSet(
new HashSet<>(