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

Reply via email to