This is an automated email from the ASF dual-hosted git repository. gaojun2048 pushed a commit to branch apache_240710_improve_event in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 98a1656c259744afb003bf9b4230ec20122b6701 Author: Eric <[email protected]> AuthorDate: Tue Jul 9 20:36:20 2024 +0800 Add EventService --- .../seatunnel/engine/server/EventService.java | 100 +++++++++++++++++++++ .../seatunnel/engine/server/SeaTunnelServer.java | 10 ++- .../engine/server/TaskExecutionService.java | 54 ++--------- .../seatunnel/engine/server/master/JobMaster.java | 9 ++ 4 files changed, 123 insertions(+), 50 deletions(-) diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/EventService.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/EventService.java new file mode 100644 index 0000000000..0c7b654b21 --- /dev/null +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/EventService.java @@ -0,0 +1,100 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.engine.server; + +import org.apache.seatunnel.shade.com.google.common.util.concurrent.ThreadFactoryBuilder; + +import org.apache.seatunnel.api.event.Event; +import org.apache.seatunnel.common.utils.RetryUtils; +import org.apache.seatunnel.engine.server.event.JobEventReportOperation; +import org.apache.seatunnel.engine.server.utils.NodeEngineUtil; + +import com.hazelcast.spi.impl.NodeEngineImpl; +import lombok.extern.slf4j.Slf4j; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + +@Slf4j +public class EventService { + private final BlockingQueue<Event> eventBuffer; + + private ExecutorService eventForwardService; + + private final NodeEngineImpl nodeEngine; + + public EventService(NodeEngineImpl nodeEngine) { + eventBuffer = new ArrayBlockingQueue<>(2048); + initEventForwardService(); + this.nodeEngine = nodeEngine; + } + + private void initEventForwardService() { + eventForwardService = + Executors.newSingleThreadExecutor( + new ThreadFactoryBuilder().setNameFormat("event-forwarder-%d").build()); + eventForwardService.submit( + () -> { + List<Event> events = new ArrayList<>(); + RetryUtils.RetryMaterial retryMaterial = + new RetryUtils.RetryMaterial(2, true, e -> true); + while (!Thread.currentThread().isInterrupted()) { + try { + events.clear(); + + Event first = eventBuffer.take(); + events.add(first); + + eventBuffer.drainTo(events, 500); + JobEventReportOperation operation = new JobEventReportOperation(events); + + RetryUtils.retryWithException( + () -> + NodeEngineUtil.sendOperationToMasterNode( + nodeEngine, operation) + .join(), + retryMaterial); + + log.debug("Event forward success, events " + events.size()); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + log.info("Event forward thread interrupted"); + } catch (Throwable t) { + log.warn("Event forward failed, discard events " + events.size(), t); + } + } + }); + } + + public void reportEvent(Event e) { + while (!eventBuffer.offer(e)) { + eventBuffer.poll(); + log.warn("Event buffer is full, discard the oldest event"); + } + } + + public void shutdownNow() { + if (eventForwardService != null) { + eventForwardService.shutdownNow(); + } + } +} diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/SeaTunnelServer.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/SeaTunnelServer.java index 765869fd03..831b5f5ea5 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/SeaTunnelServer.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/SeaTunnelServer.java @@ -76,6 +76,8 @@ public class SeaTunnelServer private volatile boolean isRunning = true; + @Getter private EventService eventService; + public SeaTunnelServer(@NonNull SeaTunnelConfig seaTunnelConfig) { this.liveOperationRegistry = new LiveOperationRegistry(); this.seaTunnelConfig = seaTunnelConfig; @@ -110,6 +112,8 @@ public class SeaTunnelServer new DefaultClassLoaderService( seaTunnelConfig.getEngineConfig().isClassloaderCacheMode()); + eventService = new EventService(nodeEngine); + if (EngineConfig.ClusterRole.MASTER_AND_WORKER.ordinal() == seaTunnelConfig.getEngineConfig().getClusterRole().ordinal()) { startWorker(); @@ -143,7 +147,7 @@ public class SeaTunnelServer private void startWorker() { taskExecutionService = new TaskExecutionService( - classLoaderService, nodeEngine, nodeEngine.getProperties()); + classLoaderService, nodeEngine, nodeEngine.getProperties(), eventService); nodeEngine.getMetricsRegistry().registerDynamicMetricsProvider(taskExecutionService); taskExecutionService.start(); getSlotService(); @@ -170,6 +174,10 @@ public class SeaTunnelServer if (coordinatorService != null) { coordinatorService.shutdown(); } + + if (eventService != null) { + eventService.shutdownNow(); + } } @Override diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/TaskExecutionService.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/TaskExecutionService.java index 19878545ed..beb51ac082 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/TaskExecutionService.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/TaskExecutionService.java @@ -20,7 +20,6 @@ package org.apache.seatunnel.engine.server; import org.apache.seatunnel.api.common.metrics.MetricTags; import org.apache.seatunnel.api.event.Event; import org.apache.seatunnel.common.utils.ExceptionUtils; -import org.apache.seatunnel.common.utils.RetryUtils; import org.apache.seatunnel.common.utils.StringFormatUtils; import org.apache.seatunnel.engine.common.Constant; import org.apache.seatunnel.engine.common.config.ConfigProvider; @@ -30,7 +29,6 @@ import org.apache.seatunnel.engine.common.exception.JobNotFoundException; import org.apache.seatunnel.engine.common.utils.PassiveCompletableFuture; import org.apache.seatunnel.engine.core.classloader.ClassLoaderService; import org.apache.seatunnel.engine.core.job.ConnectorJarIdentifier; -import org.apache.seatunnel.engine.server.event.JobEventReportOperation; import org.apache.seatunnel.engine.server.exception.TaskGroupContextNotFoundException; import org.apache.seatunnel.engine.server.execution.ExecutionState; import org.apache.seatunnel.engine.server.execution.ProgressState; @@ -49,7 +47,6 @@ import org.apache.seatunnel.engine.server.service.jar.ServerConnectorPackageClie import org.apache.seatunnel.engine.server.task.SeaTunnelTask; import org.apache.seatunnel.engine.server.task.TaskGroupImmutableInformation; import org.apache.seatunnel.engine.server.task.operation.NotifyTaskStatusOperation; -import org.apache.seatunnel.engine.server.utils.NodeEngineUtil; import org.apache.commons.collections4.CollectionUtils; @@ -73,7 +70,6 @@ import lombok.SneakyThrows; import java.io.IOException; import java.net.URL; -import java.util.ArrayList; import java.util.Collection; import java.util.HashMap; import java.util.HashSet; @@ -81,7 +77,6 @@ import java.util.List; import java.util.Map; import java.util.Set; import java.util.UUID; -import java.util.concurrent.ArrayBlockingQueue; import java.util.concurrent.BlockingQueue; import java.util.concurrent.CancellationException; import java.util.concurrent.CompletableFuture; @@ -147,13 +142,13 @@ public class TaskExecutionService implements DynamicMetricsProvider { private final ServerConnectorPackageClient serverConnectorPackageClient; - private final BlockingQueue<Event> eventBuffer; - private final ExecutorService eventForwardService; + private final EventService eventService; public TaskExecutionService( ClassLoaderService classLoaderService, NodeEngineImpl nodeEngine, - HazelcastProperties properties) { + HazelcastProperties properties, + EventService eventService) { seaTunnelConfig = ConfigProvider.locateAndGetSeaTunnelConfig(); this.hzInstanceName = nodeEngine.getHazelcastInstance().getName(); this.nodeEngine = nodeEngine; @@ -176,42 +171,7 @@ public class TaskExecutionService implements DynamicMetricsProvider { serverConnectorPackageClient = new ServerConnectorPackageClient(nodeEngine, seaTunnelConfig); - eventBuffer = new ArrayBlockingQueue<>(2048); - eventForwardService = - Executors.newSingleThreadExecutor( - new ThreadFactoryBuilder().setNameFormat("event-forwarder-%d").build()); - eventForwardService.submit( - () -> { - List<Event> events = new ArrayList<>(); - RetryUtils.RetryMaterial retryMaterial = - new RetryUtils.RetryMaterial(2, true, e -> true); - while (!Thread.currentThread().isInterrupted()) { - try { - events.clear(); - - Event first = eventBuffer.take(); - events.add(first); - - eventBuffer.drainTo(events, 500); - JobEventReportOperation operation = new JobEventReportOperation(events); - - RetryUtils.retryWithException( - () -> - NodeEngineUtil.sendOperationToMasterNode( - nodeEngine, operation) - .join(), - retryMaterial); - - logger.fine("Event forward success, events " + events.size()); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - logger.info("Event forward thread interrupted"); - } catch (Throwable t) { - logger.warning( - "Event forward failed, discard events " + events.size(), t); - } - } - }); + this.eventService = eventService; } public void start() { @@ -222,7 +182,6 @@ public class TaskExecutionService implements DynamicMetricsProvider { isRunning = false; executorService.shutdownNow(); scheduledExecutorService.shutdown(); - eventForwardService.shutdownNow(); } public TaskGroupContext getExecutionContext(TaskGroupLocation taskGroupLocation) { @@ -668,10 +627,7 @@ public class TaskExecutionService implements DynamicMetricsProvider { } public void reportEvent(Event e) { - while (!eventBuffer.offer(e)) { - eventBuffer.poll(); - logger.warning("Event buffer is full, discard the oldest event"); - } + eventService.reportEvent(e); } private final class BlockingWorker implements Runnable { diff --git a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/master/JobMaster.java b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/master/JobMaster.java index 29d8611f13..4c68d29cc4 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/master/JobMaster.java +++ b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/master/JobMaster.java @@ -384,6 +384,15 @@ public class JobMaster { } } + private void reportEventOfSaveMode( + long jobId, TablePath tablePath, int indexOfCount, long startTime, long finishedTime) { + seaTunnelServer + .getEventService() + .reportEvent( + new SaveModeFinishedEvent( + jobId, tablePath, indexOfCount, startTime, finishedTime)); + } + public void handleCheckpointError(long pipelineId, boolean neverRestore) { if (neverRestore) { this.neverNeedRestore();
