This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12023-455d3002831812819726d16376bc5ad588307bfd in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 9b4814726cd50d2c72557b8ba6e40b8c4cc16cba Author: Jast <[email protected]> AuthorDate: Sat Sep 19 14:38:41 2026 +0000 [Fix][Zeta] Skip connector jar operations on master-only nodes (#12023) Co-authored-by: zhangshenghang <[email protected]> --- .../DeleteConnectorJarInExecutionNode.java | 3 ++ .../SendConnectorJarToMemberNodeOperation.java | 3 ++ .../engine/server/ConnectorPackageServiceTest.java | 62 ++++++++++++++++++++++ 3 files changed, 68 insertions(+) diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/operation/DeleteConnectorJarInExecutionNode.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/operation/DeleteConnectorJarInExecutionNode.java index 52b211f15c..1e7d2221fd 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/operation/DeleteConnectorJarInExecutionNode.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/operation/DeleteConnectorJarInExecutionNode.java @@ -52,6 +52,9 @@ public class DeleteConnectorJarInExecutionNode extends Operation @Override public void run() throws Exception { SeaTunnelServer seaTunnelServer = getService(); + if (seaTunnelServer.getTaskExecutionService() == null) { + return; + } ServerConnectorPackageClient serverConnectorPackageClient = seaTunnelServer.getTaskExecutionService().getServerConnectorPackageClient(); serverConnectorPackageClient.deleteConnectorJar(connectorJarIdentifier); diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/operation/SendConnectorJarToMemberNodeOperation.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/operation/SendConnectorJarToMemberNodeOperation.java index 94d767cb07..53066762f0 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/operation/SendConnectorJarToMemberNodeOperation.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/task/operation/SendConnectorJarToMemberNodeOperation.java @@ -57,6 +57,9 @@ public class SendConnectorJarToMemberNodeOperation extends Operation @Override public void run() throws Exception { SeaTunnelServer seaTunnelServer = getService(); + if (seaTunnelServer.getTaskExecutionService() == null) { + return; + } ServerConnectorPackageClient serverConnectorPackageClient = seaTunnelServer.getTaskExecutionService().getServerConnectorPackageClient(); serverConnectorPackageClient.storageConnectorJarFile( diff --git a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/ConnectorPackageServiceTest.java b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/ConnectorPackageServiceTest.java index 73d3eb6115..aa6a3a576e 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/ConnectorPackageServiceTest.java +++ b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/ConnectorPackageServiceTest.java @@ -48,6 +48,9 @@ import org.apache.seatunnel.engine.core.job.JobImmutableInformation; import org.apache.seatunnel.engine.core.job.PipelineStatus; import org.apache.seatunnel.engine.core.parse.MultipleTableJobConfigParser; import org.apache.seatunnel.engine.server.service.jar.ConnectorPackageService; +import org.apache.seatunnel.engine.server.task.operation.DeleteConnectorJarInExecutionNode; +import org.apache.seatunnel.engine.server.task.operation.SendConnectorJarToMemberNodeOperation; +import org.apache.seatunnel.engine.server.utils.NodeEngineUtil; import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.Assertions; @@ -176,6 +179,65 @@ public class ConnectorPackageServiceTest { } } + @Test + public void testConnectorJarOperationsSkipMasterOnlyNode() { + SEATUNNEL_CONFIG + .getHazelcastConfig() + .setClusterName( + TestUtils.getClusterName( + "ConnectorPackageServiceTest_testConnectorJarOperationsSkipMasterOnlyNode")); + HazelcastInstanceImpl instance1 = null; + HazelcastInstanceImpl instance2 = null; + try { + instance1 = SeaTunnelServerStarter.createMasterHazelcastInstance(SEATUNNEL_CONFIG); + instance2 = SeaTunnelServerStarter.createMasterHazelcastInstance(SEATUNNEL_CONFIG); + HazelcastInstanceImpl finalInstance1 = instance1; + HazelcastInstanceImpl finalInstance2 = instance2; + await().atMost(20000, TimeUnit.MILLISECONDS) + .untilAsserted( + () -> + Assertions.assertEquals( + 2, finalInstance1.getCluster().getMembers().size())); + + SeaTunnelServer masterOnlyServer = + finalInstance2.node.getNodeEngine().getService(SeaTunnelServer.SERVICE_NAME); + Assertions.assertNull(masterOnlyServer.getTaskExecutionService()); + + ConnectorJar connectorJar = + ConnectorJar.createConnectorJar( + new byte[] {1}, + ConnectorJarType.CONNECTOR_PLUGIN_JAR, + new byte[] {1}, + "test-connector.jar"); + ConnectorJarIdentifier connectorJarIdentifier = + ConnectorJarIdentifier.of(connectorJar, "/tmp/test-connector.jar"); + + Assertions.assertDoesNotThrow( + () -> + NodeEngineUtil.sendOperationToMemberNode( + finalInstance1.node.getNodeEngine(), + new SendConnectorJarToMemberNodeOperation( + connectorJar, connectorJarIdentifier), + finalInstance2.node.getThisAddress()) + .join()); + Assertions.assertDoesNotThrow( + () -> + NodeEngineUtil.sendOperationToMemberNode( + finalInstance1.node.getNodeEngine(), + new DeleteConnectorJarInExecutionNode( + connectorJarIdentifier), + finalInstance2.node.getThisAddress()) + .join()); + } finally { + if (instance2 != null && instance2.node.isRunning()) { + instance2.shutdown(); + } + if (instance1 != null && instance1.node.isRunning()) { + instance1.shutdown(); + } + } + } + @Test @Disabled("disabled because we can not know") public void testRestoreWhenMasterNodeSwitch() throws InterruptedException, IOException {
