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<>(
