This is an automated email from the ASF dual-hosted git repository.

jonyang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git


The following commit(s) were added to refs/heads/master by this push:
     new d00049205 [ISSUE #3126]Thread#sleep replace with ThreadUtils#sleep 
(#3127)
d00049205 is described below

commit d000492052ee36c31087b2f578bc9f64a1f6b415
Author: mxsm <[email protected]>
AuthorDate: Sun Feb 12 22:51:59 2023 +0800

    [ISSUE #3126]Thread#sleep replace with ThreadUtils#sleep (#3127)
    
    * [ISSUE #3126]Thread#sleep replace with ThreadUtils#sleep
    
    * sleep(long timeout) add @deprecated annotation and polish code
---
 .../org/apache/eventmesh/common/utils/ThreadUtils.java   | 13 +++++++++----
 .../eventmesh/common/file/WatchFileManagerTest.java      |  5 ++++-
 .../rabbitmq/consumer/RabbitmqConsumerTest.java          |  3 ++-
 .../rabbitmq/producer/RabbitmqProducerTest.java          |  3 ++-
 .../impl/consumer/ConsumeMessageConcurrentlyService.java |  7 ++-----
 .../standalone/broker/task/HistoryMessageClearTask.java  |  3 ++-
 .../connector/standalone/broker/task/SubScribeTask.java  |  4 +++-
 .../pub/cloudevents/CloudEventsBatchPublishInstance.java |  4 +++-
 .../grpc/pub/cloudevents/CloudEventsPublishInstance.java |  7 +++++--
 .../grpc/pub/cloudevents/CloudEventsRequestInstance.java |  7 +++++--
 .../grpc/pub/eventmeshmessage/AsyncPublishBroadcast.java |  7 +++++--
 .../grpc/pub/eventmeshmessage/AsyncPublishInstance.java  |  7 +++++--
 .../grpc/pub/eventmeshmessage/BatchPublishInstance.java  |  5 ++++-
 .../grpc/pub/eventmeshmessage/RequestReplyInstance.java  |  7 +++++--
 .../eventmeshmessage/WorkflowAsyncPublishInstance.java   |  5 ++++-
 .../eventmesh/grpc/sub/CloudEventsAsyncSubscribe.java    |  4 +++-
 .../eventmesh/grpc/sub/CloudEventsSubscribeReply.java    |  4 +++-
 .../eventmesh/grpc/sub/EventmeshAsyncSubscribe.java      |  5 ++++-
 .../eventmesh/grpc/sub/EventmeshSubscribeBroadcast.java  |  5 ++++-
 .../eventmesh/grpc/sub/EventmeshSubscribeReply.java      |  7 +++++--
 .../grpc/sub/WorkflowExpressAsyncSubscribe.java          |  5 ++++-
 .../eventmesh/grpc/sub/WorkflowOrderAsyncSubscribe.java  |  5 ++++-
 .../grpc/sub/WorkflowPaymentAsyncSubscribe.java          |  5 ++++-
 .../http/demo/pub/cloudevents/AsyncPublishInstance.java  |  5 ++++-
 .../demo/pub/eventmeshmessage/AsyncPublishInstance.java  |  5 ++++-
 .../pub/eventmeshmessage/AsyncSyncRequestInstance.java   |  8 ++++++--
 .../demo/pub/eventmeshmessage/SyncRequestInstance.java   |  6 +++++-
 .../eventmesh/tcp/demo/pub/cloudevents/AsyncPublish.java |  6 ++++--
 .../tcp/demo/pub/eventmeshmessage/AsyncPublish.java      |  7 +++++--
 .../demo/pub/eventmeshmessage/AsyncPublishBroadcast.java |  5 ++++-
 .../eventmesh/runtime/boot/EventMeshTCPServer.java       |  9 +++------
 .../core/protocol/grpc/consumer/EventMeshConsumer.java   |  6 ++++--
 .../tcp/client/group/ClientSessionGroupMapping.java      | 16 +++++-----------
 .../tcp/client/rebalance/EventmeshRebalanceImpl.java     |  8 +++-----
 .../eventmesh/client/http/demo/AsyncPublishInstance.java |  7 +++++--
 .../client/http/demo/AsyncSyncRequestInstance.java       |  6 ++++--
 .../eventmesh/client/http/demo/SyncRequestInstance.java  |  4 +++-
 37 files changed, 150 insertions(+), 75 deletions(-)

diff --git 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/utils/ThreadUtils.java
 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/utils/ThreadUtils.java
index 28a8121ca..b68050e17 100644
--- 
a/eventmesh-common/src/main/java/org/apache/eventmesh/common/utils/ThreadUtils.java
+++ 
b/eventmesh-common/src/main/java/org/apache/eventmesh/common/utils/ThreadUtils.java
@@ -44,21 +44,26 @@ public class ThreadUtils {
         randomPause(1, max);
     }
 
+    @Deprecated
     public static void sleep(long timeout) {
         sleep(timeout, TimeUnit.MILLISECONDS);
     }
 
     public static void sleep(long timeout, TimeUnit timeUnit) {
-        if (null == timeUnit) {
-            return;
-        }
         try {
-            timeUnit.sleep(timeout);
+            sleepWithThrowException(timeout, timeUnit);
         } catch (InterruptedException ignore) {
             //ignore
         }
     }
 
+    public static void sleepWithThrowException(long timeout, TimeUnit 
timeUnit) throws InterruptedException {
+        if (null == timeUnit) {
+            return;
+        }
+        timeUnit.sleep(timeout);
+    }
+
     /**
      * get current process id.
      *
diff --git 
a/eventmesh-common/src/test/java/org/apache/eventmesh/common/file/WatchFileManagerTest.java
 
b/eventmesh-common/src/test/java/org/apache/eventmesh/common/file/WatchFileManagerTest.java
index 614138a4a..63448220c 100644
--- 
a/eventmesh-common/src/test/java/org/apache/eventmesh/common/file/WatchFileManagerTest.java
+++ 
b/eventmesh-common/src/test/java/org/apache/eventmesh/common/file/WatchFileManagerTest.java
@@ -17,12 +17,15 @@
 
 package org.apache.eventmesh.common.file;
 
+import org.apache.eventmesh.common.utils.ThreadUtils;
+
 import java.io.BufferedReader;
 import java.io.File;
 import java.io.FileReader;
 import java.io.FileWriter;
 import java.io.IOException;
 import java.util.Properties;
+import java.util.concurrent.TimeUnit;
 
 import org.junit.Assert;
 import org.junit.Test;
@@ -53,6 +56,6 @@ public class WatchFileManagerTest {
         FileWriter fw = new FileWriter(file);
         properties.store(fw, "newAdd");
 
-        Thread.sleep(500);
+        ThreadUtils.sleep(500, TimeUnit.MILLISECONDS);
     }
 }
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/test/java/org/apache/eventmesh/connector/rabbitmq/consumer/RabbitmqConsumerTest.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/test/java/org/apache/eventmesh/connector/rabbitmq/consumer/RabbitmqConsumerTest.java
index 1e5a3969a..ebf9e25fa 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/test/java/org/apache/eventmesh/connector/rabbitmq/consumer/RabbitmqConsumerTest.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/test/java/org/apache/eventmesh/connector/rabbitmq/consumer/RabbitmqConsumerTest.java
@@ -21,6 +21,7 @@ import org.apache.eventmesh.api.EventMeshAction;
 import org.apache.eventmesh.api.SendCallback;
 import org.apache.eventmesh.api.SendResult;
 import org.apache.eventmesh.api.exception.OnExceptionContext;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.connector.rabbitmq.RabbitmqServer;
 
 import java.net.URI;
@@ -60,7 +61,7 @@ public class RabbitmqConsumerTest extends RabbitmqServer {
 
         rabbitmqConsumer.subscribe("topic");
 
-        Thread.sleep(1000);
+        ThreadUtils.sleep(1, TimeUnit.SECONDS);
         for (int i = 0; i < expectedCount; i++) {
             CloudEvent cloudEvent = CloudEventBuilder.v1()
                     .withId(String.valueOf(i))
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/test/java/org/apache/eventmesh/connector/rabbitmq/producer/RabbitmqProducerTest.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/test/java/org/apache/eventmesh/connector/rabbitmq/producer/RabbitmqProducerTest.java
index bf39e2cc5..567e74720 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/test/java/org/apache/eventmesh/connector/rabbitmq/producer/RabbitmqProducerTest.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rabbitmq/src/test/java/org/apache/eventmesh/connector/rabbitmq/producer/RabbitmqProducerTest.java
@@ -21,6 +21,7 @@ import org.apache.eventmesh.api.EventMeshAction;
 import org.apache.eventmesh.api.SendCallback;
 import org.apache.eventmesh.api.SendResult;
 import org.apache.eventmesh.api.exception.OnExceptionContext;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.connector.rabbitmq.RabbitmqServer;
 
 import java.net.URI;
@@ -60,7 +61,7 @@ public class RabbitmqProducerTest extends RabbitmqServer {
 
         rabbitmqConsumer.subscribe("topic");
 
-        Thread.sleep(1000);
+        ThreadUtils.sleep(1, TimeUnit.SECONDS);
         for (int i = 0; i < expectedCount; i++) {
             CloudEvent cloudEvent = CloudEventBuilder.v1()
                     .withId(String.valueOf(i))
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
index 5cfe1bc73..e2119ec2e 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
@@ -18,6 +18,7 @@
 package org.apache.rocketmq.client.impl.consumer;
 
 import org.apache.eventmesh.common.EventMeshThreadFactory;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import 
org.apache.eventmesh.connector.rocketmq.patch.EventMeshConsumeConcurrentlyContext;
 
 import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
@@ -371,11 +372,7 @@ public class ConsumeMessageConcurrentlyService implements 
ConsumeMessageService
                 boolean success = false;
                 for (int i = 0; i < times; i++) {
                     try {
-                        try {
-                            Thread.sleep(1000);
-                        } catch (InterruptedException e) {
-                            //ignore
-                        }
+                        ThreadUtils.sleep(1, TimeUnit.SECONDS);
                         
ConsumeMessageConcurrentlyService.this.consumeExecutor.submit(consumeRequest);
                         success = true;
                         break;
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/HistoryMessageClearTask.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/HistoryMessageClearTask.java
index 5499f331f..f8aece4b3 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/HistoryMessageClearTask.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/HistoryMessageClearTask.java
@@ -17,6 +17,7 @@
 
 package org.apache.eventmesh.connector.standalone.broker.task;
 
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.connector.standalone.broker.MessageQueue;
 import org.apache.eventmesh.connector.standalone.broker.model.MessageEntity;
 import org.apache.eventmesh.connector.standalone.broker.model.TopicMetadata;
@@ -60,7 +61,7 @@ public class HistoryMessageClearTask implements Runnable {
                 }
             });
             try {
-                Thread.sleep(TimeUnit.SECONDS.toMillis(1));
+                ThreadUtils.sleepWithThrowException(1, TimeUnit.SECONDS);
             } catch (InterruptedException e) {
                 logger.error("Thread is interrupted, thread name: {}", 
Thread.currentThread().getName(), e);
                 Thread.currentThread().interrupt();
diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
index aec2adda9..5e9dcfd70 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-standalone/src/main/java/org/apache/eventmesh/connector/standalone/broker/task/SubScribeTask.java
@@ -20,8 +20,10 @@ package 
org.apache.eventmesh.connector.standalone.broker.task;
 import org.apache.eventmesh.api.EventListener;
 import org.apache.eventmesh.api.EventMeshAction;
 import org.apache.eventmesh.api.EventMeshAsyncConsumeContext;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.connector.standalone.broker.StandaloneBroker;
 
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 
 import org.slf4j.Logger;
@@ -101,7 +103,7 @@ public class SubScribeTask implements Runnable {
                     ex);
             }
             try {
-                Thread.sleep(1000);
+                ThreadUtils.sleepWithThrowException(1, TimeUnit.SECONDS);
             } catch (InterruptedException e) {
                 logger.error("Thread is interrupted, topic: {}, offset: {} 
thread name: {}",
                     topicName, offset == null ? null : offset.get(), 
Thread.currentThread().getName(), e);
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsBatchPublishInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsBatchPublishInstance.java
index 3334cfa76..f18a7c201 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsBatchPublishInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsBatchPublishInstance.java
@@ -19,12 +19,14 @@ package org.apache.eventmesh.grpc.pub.cloudevents;
 
 import org.apache.eventmesh.client.grpc.producer.EventMeshGrpcProducer;
 import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
 
 import io.cloudevents.CloudEvent;
 
@@ -46,7 +48,7 @@ public class CloudEventsBatchPublishInstance extends 
GrpcAbstractDemo {
                 cloudEventList.add(buildCloudEvent(content));
             }
             eventMeshGrpcProducer.publish(cloudEventList);
-            Thread.sleep(10_000);
+            ThreadUtils.sleep(10, TimeUnit.SECONDS);
         }
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsPublishInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsPublishInstance.java
index 0bcc0e244..6d4382d10 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsPublishInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsPublishInstance.java
@@ -19,10 +19,13 @@ package org.apache.eventmesh.grpc.pub.cloudevents;
 
 import org.apache.eventmesh.client.grpc.producer.EventMeshGrpcProducer;
 import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -42,10 +45,10 @@ public class CloudEventsPublishInstance extends 
GrpcAbstractDemo {
 
             for (int i = 0; i < MESSAGE_SIZE; i++) {
                 eventMeshGrpcProducer.publish(buildCloudEvent(content));
-                Thread.sleep(1_000);
+                ThreadUtils.sleep(1, TimeUnit.SECONDS);
             }
 
-            Thread.sleep(30_000);
+            ThreadUtils.sleep(30, TimeUnit.SECONDS);
         }
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsRequestInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsRequestInstance.java
index 7f862c35d..04aab6fab 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsRequestInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/cloudevents/CloudEventsRequestInstance.java
@@ -20,10 +20,13 @@ package org.apache.eventmesh.grpc.pub.cloudevents;
 import org.apache.eventmesh.client.grpc.producer.EventMeshGrpcProducer;
 import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
 import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -41,10 +44,10 @@ public class CloudEventsRequestInstance extends 
GrpcAbstractDemo {
 
             for (int i = 0; i < MESSAGE_SIZE; i++) {
                 eventMeshGrpcProducer.requestReply(buildCloudEvent(content), 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
-                Thread.sleep(1_000);
+                ThreadUtils.sleep(1, TimeUnit.SECONDS);
             }
 
-            Thread.sleep(30_000);
+            ThreadUtils.sleep(30, TimeUnit.SECONDS);
         }
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/AsyncPublishBroadcast.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/AsyncPublishBroadcast.java
index 63be2397f..9b8ebf4d9 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/AsyncPublishBroadcast.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/AsyncPublishBroadcast.java
@@ -19,10 +19,13 @@ package org.apache.eventmesh.grpc.pub.eventmeshmessage;
 
 import org.apache.eventmesh.client.grpc.producer.EventMeshGrpcProducer;
 import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -41,9 +44,9 @@ public class AsyncPublishBroadcast extends GrpcAbstractDemo {
 
             for (int i = 0; i < MESSAGE_SIZE; i++) {
                 eventMeshGrpcProducer.publish(buildEventMeshMessage(content));
-                Thread.sleep(1_000);
+                ThreadUtils.sleep(1, TimeUnit.SECONDS);
             }
-            Thread.sleep(30_000);
+            ThreadUtils.sleep(30, TimeUnit.SECONDS);
         }
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/AsyncPublishInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/AsyncPublishInstance.java
index 85a7e4e72..65916040b 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/AsyncPublishInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/AsyncPublishInstance.java
@@ -19,10 +19,13 @@ package org.apache.eventmesh.grpc.pub.eventmeshmessage;
 
 import org.apache.eventmesh.client.grpc.producer.EventMeshGrpcProducer;
 import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -41,9 +44,9 @@ public class AsyncPublishInstance extends GrpcAbstractDemo {
 
             for (int i = 0; i < MESSAGE_SIZE; i++) {
                 buildEventMeshMessage(content);
-                Thread.sleep(1_000);
+                ThreadUtils.sleep(1, TimeUnit.SECONDS);
             }
-            Thread.sleep(30_000);
+            ThreadUtils.sleep(30, TimeUnit.SECONDS);
         }
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/BatchPublishInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/BatchPublishInstance.java
index e847ecba5..ea797c76d 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/BatchPublishInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/BatchPublishInstance.java
@@ -20,12 +20,15 @@ package org.apache.eventmesh.grpc.pub.eventmeshmessage;
 import org.apache.eventmesh.client.grpc.producer.EventMeshGrpcProducer;
 import org.apache.eventmesh.common.EventMeshMessage;
 import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -46,7 +49,7 @@ public class BatchPublishInstance extends GrpcAbstractDemo {
             }
 
             eventMeshGrpcProducer.publish(messageList);
-            Thread.sleep(10_000);
+            ThreadUtils.sleep(10, TimeUnit.SECONDS);
         }
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/RequestReplyInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/RequestReplyInstance.java
index 9e310fe10..3f7c050c9 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/RequestReplyInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/RequestReplyInstance.java
@@ -20,10 +20,13 @@ package org.apache.eventmesh.grpc.pub.eventmeshmessage;
 import org.apache.eventmesh.client.grpc.producer.EventMeshGrpcProducer;
 import org.apache.eventmesh.client.tcp.common.EventMeshCommon;
 import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -44,9 +47,9 @@ public class RequestReplyInstance extends GrpcAbstractDemo {
             for (int i = 0; i < MESSAGE_SIZE; i++) {
                 
eventMeshGrpcProducer.requestReply(buildEventMeshMessage(content),
                         EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
-                Thread.sleep(1_000);
+                ThreadUtils.sleep(1, TimeUnit.SECONDS);
             }
-            Thread.sleep(30_000);
+            ThreadUtils.sleep(30, TimeUnit.SECONDS);
 
         }
     }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/WorkflowAsyncPublishInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/WorkflowAsyncPublishInstance.java
index 0989d7956..08bfbbb70 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/WorkflowAsyncPublishInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/pub/eventmeshmessage/WorkflowAsyncPublishInstance.java
@@ -24,6 +24,7 @@ import 
org.apache.eventmesh.client.workflow.config.EventMeshWorkflowClientConfig
 import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteRequest;
 import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteResponse;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 import org.apache.eventmesh.selector.NacosSelector;
 import org.apache.eventmesh.util.Utils;
@@ -31,6 +32,8 @@ import org.apache.eventmesh.util.Utils;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Properties;
+import java.util.concurrent.TimeUnit;
+
 
 import com.alibaba.nacos.shaded.com.google.gson.Gson;
 
@@ -68,7 +71,7 @@ public class WorkflowAsyncPublishInstance extends 
GrpcAbstractDemo {
                 log.info("received response: {}", response.toString());
             }
 
-            Thread.sleep(60_000);
+            ThreadUtils.sleep(1, TimeUnit.MINUTES);
 
         }
     }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsAsyncSubscribe.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsAsyncSubscribe.java
index 5e26c6233..e1bfdb001 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsAsyncSubscribe.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsAsyncSubscribe.java
@@ -24,11 +24,13 @@ import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.SubscriptionItem;
 import org.apache.eventmesh.common.protocol.SubscriptionMode;
 import org.apache.eventmesh.common.protocol.SubscriptionType;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.io.IOException;
 import java.util.Collections;
 import java.util.Optional;
+import java.util.concurrent.TimeUnit;
 
 import io.cloudevents.CloudEvent;
 
@@ -52,7 +54,7 @@ public class CloudEventsAsyncSubscribe extends 
GrpcAbstractDemo implements Recei
 
             
eventMeshGrpcConsumer.subscribe(Collections.singletonList(subscriptionItem));
 
-            Thread.sleep(60_000);
+            ThreadUtils.sleep(1, TimeUnit.MINUTES);
             
eventMeshGrpcConsumer.unsubscribe(Collections.singletonList(subscriptionItem));
         }
     }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsSubscribeReply.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsSubscribeReply.java
index 5a4e8606e..048e86367 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsSubscribeReply.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/CloudEventsSubscribeReply.java
@@ -24,11 +24,13 @@ import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.SubscriptionItem;
 import org.apache.eventmesh.common.protocol.SubscriptionMode;
 import org.apache.eventmesh.common.protocol.SubscriptionType;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.io.IOException;
 import java.util.Collections;
 import java.util.Optional;
+import java.util.concurrent.TimeUnit;
 
 import io.cloudevents.CloudEvent;
 
@@ -53,7 +55,7 @@ public class CloudEventsSubscribeReply extends 
GrpcAbstractDemo implements Recei
 
             
eventMeshGrpcConsumer.subscribe(Collections.singletonList(subscriptionItem));
 
-            Thread.sleep(60_000);
+            ThreadUtils.sleep(1, TimeUnit.MINUTES);
             
eventMeshGrpcConsumer.unsubscribe(Collections.singletonList(subscriptionItem));
         }
     }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshAsyncSubscribe.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshAsyncSubscribe.java
index 945b49454..77660799c 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshAsyncSubscribe.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshAsyncSubscribe.java
@@ -25,11 +25,14 @@ import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.SubscriptionItem;
 import org.apache.eventmesh.common.protocol.SubscriptionMode;
 import org.apache.eventmesh.common.protocol.SubscriptionType;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.io.IOException;
 import java.util.Collections;
 import java.util.Optional;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -51,7 +54,7 @@ public class EventmeshAsyncSubscribe extends GrpcAbstractDemo 
implements Receive
 
             
eventMeshGrpcConsumer.subscribe(Collections.singletonList(subscriptionItem));
 
-            Thread.sleep(60_000);
+            ThreadUtils.sleep(1, TimeUnit.MINUTES);
             
eventMeshGrpcConsumer.unsubscribe(Collections.singletonList(subscriptionItem));
         }
     }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshSubscribeBroadcast.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshSubscribeBroadcast.java
index 4e847ce82..8c264a83e 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshSubscribeBroadcast.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshSubscribeBroadcast.java
@@ -25,11 +25,14 @@ import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.SubscriptionItem;
 import org.apache.eventmesh.common.protocol.SubscriptionMode;
 import org.apache.eventmesh.common.protocol.SubscriptionType;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.io.IOException;
 import java.util.Collections;
 import java.util.Optional;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -52,7 +55,7 @@ public class EventmeshSubscribeBroadcast extends 
GrpcAbstractDemo implements Rec
 
             
eventMeshGrpcConsumer.subscribe(Collections.singletonList(subscriptionItem));
 
-            Thread.sleep(60_000);
+            ThreadUtils.sleep(1, TimeUnit.MINUTES);
             
eventMeshGrpcConsumer.unsubscribe(Collections.singletonList(subscriptionItem));
         }
     }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshSubscribeReply.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshSubscribeReply.java
index 439952467..85cff545f 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshSubscribeReply.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/EventmeshSubscribeReply.java
@@ -25,11 +25,14 @@ import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.SubscriptionItem;
 import org.apache.eventmesh.common.protocol.SubscriptionMode;
 import org.apache.eventmesh.common.protocol.SubscriptionType;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 
 import java.io.IOException;
 import java.util.Collections;
 import java.util.Optional;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -51,7 +54,7 @@ public class EventmeshSubscribeReply extends GrpcAbstractDemo 
implements Receive
 
             
eventMeshGrpcConsumer.subscribe(Collections.singletonList(subscriptionItem));
 
-            Thread.sleep(60_000);
+            ThreadUtils.sleep(1, TimeUnit.MINUTES);
             
eventMeshGrpcConsumer.unsubscribe(Collections.singletonList(subscriptionItem));
         }
     }
@@ -72,4 +75,4 @@ public class EventmeshSubscribeReply extends GrpcAbstractDemo 
implements Receive
     public String getProtocolType() {
         return EventMeshCommon.EM_MESSAGE_PROTOCOL_NAME;
     }
-}
\ No newline at end of file
+}
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowExpressAsyncSubscribe.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowExpressAsyncSubscribe.java
index d5e83cc88..64db63259 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowExpressAsyncSubscribe.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowExpressAsyncSubscribe.java
@@ -29,6 +29,7 @@ import org.apache.eventmesh.common.EventMeshMessage;
 import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteRequest;
 import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteResponse;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 import org.apache.eventmesh.selector.NacosSelector;
 import org.apache.eventmesh.util.Utils;
@@ -36,6 +37,8 @@ import org.apache.eventmesh.util.Utils;
 import java.util.Map;
 import java.util.Optional;
 import java.util.Properties;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -71,7 +74,7 @@ public class WorkflowExpressAsyncSubscribe extends 
GrpcAbstractDemo implements R
                     .serverName(workflowServerName).build();
             workflowClient = new 
EventMeshWorkflowClient(eventMeshWorkflowClientConfig);
 
-            Thread.sleep(60_000_000);
+            ThreadUtils.sleep(60_000, TimeUnit.SECONDS);
             eventMeshCatalogClient.destroy();
         }
     }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowOrderAsyncSubscribe.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowOrderAsyncSubscribe.java
index 306effdb8..12f679dd0 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowOrderAsyncSubscribe.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowOrderAsyncSubscribe.java
@@ -29,6 +29,7 @@ import org.apache.eventmesh.common.EventMeshMessage;
 import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteRequest;
 import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteResponse;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 import org.apache.eventmesh.selector.NacosSelector;
 import org.apache.eventmesh.util.Utils;
@@ -36,6 +37,8 @@ import org.apache.eventmesh.util.Utils;
 import java.util.Map;
 import java.util.Optional;
 import java.util.Properties;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -71,7 +74,7 @@ public class WorkflowOrderAsyncSubscribe extends 
GrpcAbstractDemo implements Rec
                     .serverName(workflowServerName).build();
             workflowClient = new 
EventMeshWorkflowClient(eventMeshWorkflowClientConfig);
 
-            Thread.sleep(60_000_000);
+            ThreadUtils.sleep(60_000, TimeUnit.SECONDS);
             eventMeshCatalogClient.destroy();
         }
     }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowPaymentAsyncSubscribe.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowPaymentAsyncSubscribe.java
index db3904579..15b905004 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowPaymentAsyncSubscribe.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/grpc/sub/WorkflowPaymentAsyncSubscribe.java
@@ -29,6 +29,7 @@ import org.apache.eventmesh.common.EventMeshMessage;
 import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteRequest;
 import org.apache.eventmesh.common.protocol.workflow.protos.ExecuteResponse;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.grpc.GrpcAbstractDemo;
 import org.apache.eventmesh.selector.NacosSelector;
 import org.apache.eventmesh.util.Utils;
@@ -36,6 +37,8 @@ import org.apache.eventmesh.util.Utils;
 import java.util.Map;
 import java.util.Optional;
 import java.util.Properties;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -71,7 +74,7 @@ public class WorkflowPaymentAsyncSubscribe extends 
GrpcAbstractDemo implements R
                     .serverName(workflowServerName).build();
             workflowClient = new 
EventMeshWorkflowClient(eventMeshWorkflowClientConfig);
 
-            Thread.sleep(60_000_000);
+            ThreadUtils.sleep(60_000, TimeUnit.SECONDS);
             eventMeshCatalogClient.destroy();
         }
     }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/cloudevents/AsyncPublishInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/cloudevents/AsyncPublishInstance.java
index 0957baa4d..4d2b92f25 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/cloudevents/AsyncPublishInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/cloudevents/AsyncPublishInstance.java
@@ -19,10 +19,13 @@ package org.apache.eventmesh.http.demo.pub.cloudevents;
 
 import org.apache.eventmesh.client.http.producer.EventMeshHttpProducer;
 import org.apache.eventmesh.common.ExampleConstants;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.http.demo.HttpAbstractDemo;
 
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -41,7 +44,7 @@ public class AsyncPublishInstance extends HttpAbstractDemo {
                 eventMeshHttpProducer.publish(buildCloudEvent(content));
                 log.info("publish event success content: {}", content);
             }
-            Thread.sleep(30_000);
+            ThreadUtils.sleep(30, TimeUnit.SECONDS);
         }
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncPublishInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncPublishInstance.java
index 52b67d002..870174430 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncPublishInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncPublishInstance.java
@@ -23,10 +23,13 @@ import org.apache.eventmesh.common.EventMeshMessage;
 import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.utils.JsonUtils;
 import org.apache.eventmesh.common.utils.RandomStringUtils;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.http.demo.HttpAbstractDemo;
 
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -52,7 +55,7 @@ public class AsyncPublishInstance extends HttpAbstractDemo {
                         .addProp(Constants.EVENTMESH_MESSAGE_CONST_TTL, 
String.valueOf(4 * 1000));
                 eventMeshHttpProducer.publish(eventMeshMessage);
             }
-            Thread.sleep(30_000);
+            ThreadUtils.sleep(30, TimeUnit.SECONDS);
         }
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncSyncRequestInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncSyncRequestInstance.java
index d3691d33a..7d30ce02b 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncSyncRequestInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/AsyncSyncRequestInstance.java
@@ -22,8 +22,12 @@ import org.apache.eventmesh.client.http.producer.RRCallback;
 import org.apache.eventmesh.common.EventMeshMessage;
 import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.utils.RandomStringUtils;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.http.demo.HttpAbstractDemo;
 
+import java.util.concurrent.TimeUnit;
+
+
 import lombok.extern.slf4j.Slf4j;
 
 
@@ -57,11 +61,11 @@ public class AsyncSyncRequestInstance extends 
HttpAbstractDemo {
                 }
             }, 3_000);
 
-            Thread.sleep(2_000);
+            ThreadUtils.sleep(2, TimeUnit.SECONDS);
         } catch (Exception e) {
             log.error("async send msg failed", e);
         }
-        Thread.sleep(30_000);
+        ThreadUtils.sleep(30, TimeUnit.SECONDS);
 
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/SyncRequestInstance.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/SyncRequestInstance.java
index ab34c2e9a..4a8f7928c 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/SyncRequestInstance.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/http/demo/pub/eventmeshmessage/SyncRequestInstance.java
@@ -21,8 +21,12 @@ import 
org.apache.eventmesh.client.http.producer.EventMeshHttpProducer;
 import org.apache.eventmesh.common.EventMeshMessage;
 import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.utils.RandomStringUtils;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.http.demo.HttpAbstractDemo;
 
+import java.util.concurrent.TimeUnit;
+
+
 import lombok.extern.slf4j.Slf4j;
 
 @Slf4j
@@ -52,7 +56,7 @@ public class SyncRequestInstance extends HttpAbstractDemo {
             log.error("send msg failed, ", e);
         }
 
-        Thread.sleep(30_000);
+        ThreadUtils.sleep(30, TimeUnit.SECONDS);
 
     }
 }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/cloudevents/AsyncPublish.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/cloudevents/AsyncPublish.java
index 237ae259b..d6bef6fa3 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/cloudevents/AsyncPublish.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/cloudevents/AsyncPublish.java
@@ -23,10 +23,12 @@ import 
org.apache.eventmesh.client.tcp.common.EventMeshCommon;
 import org.apache.eventmesh.client.tcp.conf.EventMeshTCPClientConfig;
 import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.tcp.UserAgent;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.tcp.common.EventMeshTestUtils;
 import org.apache.eventmesh.util.Utils;
 
 import java.util.Properties;
+import java.util.concurrent.TimeUnit;
 
 import io.cloudevents.CloudEvent;
 
@@ -57,9 +59,9 @@ public class AsyncPublish {
                 }
                 client.publish(event, EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
 
-                Thread.sleep(1_000);
+                ThreadUtils.sleep(1, TimeUnit.SECONDS);
             }
-            Thread.sleep(2_000);
+            ThreadUtils.sleep(2, TimeUnit.SECONDS);
         } catch (Exception e) {
             log.error("AsyncPublish failed", e);
         }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublish.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublish.java
index 866538589..d44be7272 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublish.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublish.java
@@ -24,10 +24,13 @@ import 
org.apache.eventmesh.client.tcp.conf.EventMeshTCPClientConfig;
 import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
 import org.apache.eventmesh.common.protocol.tcp.UserAgent;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.tcp.common.EventMeshTestUtils;
 import org.apache.eventmesh.util.Utils;
 
 import java.util.Properties;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -57,9 +60,9 @@ public class AsyncPublish {
                 }
                 client.publish(eventMeshMessage, 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
 
-                Thread.sleep(1_000);
+                ThreadUtils.sleep(1, TimeUnit.SECONDS);
             }
-            Thread.sleep(2_000);
+            ThreadUtils.sleep(2, TimeUnit.SECONDS);
         } catch (Exception e) {
             log.error("AsyncPublish failed", e);
         }
diff --git 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublishBroadcast.java
 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublishBroadcast.java
index e22c56089..97846b931 100644
--- 
a/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublishBroadcast.java
+++ 
b/eventmesh-examples/src/main/java/org/apache/eventmesh/tcp/demo/pub/eventmeshmessage/AsyncPublishBroadcast.java
@@ -24,10 +24,13 @@ import 
org.apache.eventmesh.client.tcp.conf.EventMeshTCPClientConfig;
 import org.apache.eventmesh.common.ExampleConstants;
 import org.apache.eventmesh.common.protocol.tcp.EventMeshMessage;
 import org.apache.eventmesh.common.protocol.tcp.UserAgent;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.tcp.common.EventMeshTestUtils;
 import org.apache.eventmesh.util.Utils;
 
 import java.util.Properties;
+import java.util.concurrent.TimeUnit;
+
 
 import lombok.extern.slf4j.Slf4j;
 
@@ -56,7 +59,7 @@ public class AsyncPublishBroadcast {
             }
             client.broadcast(eventMeshMessage, 
EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
 
-            Thread.sleep(2_000);
+            ThreadUtils.sleep(2, TimeUnit.SECONDS);
 
         } catch (Exception e) {
             log.error("AsyncPublishBroadcast failed", e);
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTCPServer.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTCPServer.java
index 034662636..91a9fb84f 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTCPServer.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/boot/EventMeshTCPServer.java
@@ -25,6 +25,7 @@ import 
org.apache.eventmesh.common.exception.EventMeshException;
 import org.apache.eventmesh.common.protocol.tcp.codec.Codec;
 import org.apache.eventmesh.common.utils.ConfigurationContextUtil;
 import org.apache.eventmesh.common.utils.IPUtils;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.metrics.api.MetricsPluginFactory;
 import org.apache.eventmesh.metrics.api.MetricsRegistry;
 import org.apache.eventmesh.runtime.configuration.EventMeshTCPConfiguration;
@@ -45,6 +46,7 @@ import java.util.Optional;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.LinkedBlockingQueue;
 import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
 
 import org.assertj.core.util.Lists;
 import org.slf4j.Logger;
@@ -269,12 +271,7 @@ public class EventMeshTCPServer extends 
AbstractRemotingServer {
         }
 
         clientSessionGroupMapping.shutdown();
-        try {
-            Thread.sleep(40 * 1000);
-        } catch (InterruptedException e) {
-            LOGGER.error("interruptedException occurred while sleeping", e);
-        }
-
+        ThreadUtils.sleep(40, TimeUnit.SECONDS);
         globalTrafficShapingHandler.release();
 
         if (this.getIoGroup() != null) {
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/EventMeshConsumer.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/EventMeshConsumer.java
index dc08f86f3..b2d5b6efa 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/EventMeshConsumer.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/grpc/consumer/EventMeshConsumer.java
@@ -33,6 +33,7 @@ import org.apache.eventmesh.api.exception.OnExceptionContext;
 import org.apache.eventmesh.common.Constants;
 import org.apache.eventmesh.common.protocol.grpc.protos.Subscription;
 import 
org.apache.eventmesh.common.protocol.grpc.protos.Subscription.SubscriptionItem.SubscriptionMode;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.runtime.boot.EventMeshGrpcServer;
 import org.apache.eventmesh.runtime.common.ServiceState;
 import org.apache.eventmesh.runtime.configuration.EventMeshGrpcConfiguration;
@@ -53,6 +54,7 @@ import java.util.Map;
 import java.util.Optional;
 import java.util.Properties;
 import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.TimeUnit;
 
 import io.cloudevents.CloudEvent;
 import io.cloudevents.core.builder.CloudEventBuilder;
@@ -274,9 +276,9 @@ public class EventMeshConsumer {
                     return;
                 } else {
                     // can not handle the message due to the capacity limit is 
reached
-                    // wait for sometime and send this message back to mq and 
consume again
+                    // wait for some time and send this message back to mq and 
consume again
                     try {
-                        Thread.sleep(5000);
+                        ThreadUtils.sleep(5, TimeUnit.SECONDS);
                         sendMessageBack(consumerGroup, event, uniqueId, 
bizSeqNo);
                     } catch (Exception ignored) {
                         // ignore exception
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientSessionGroupMapping.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientSessionGroupMapping.java
index b9c6547f8..0034b3d55 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientSessionGroupMapping.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientSessionGroupMapping.java
@@ -322,7 +322,7 @@ public class ClientSessionGroupMapping {
     }
 
     private void cleanClientGroupWrapperCommon(ClientGroupWrapper 
clientGroupWrapper) throws Exception {
-        
+
         if 
(CollectionUtils.isEmpty(clientGroupWrapper.getGroupConsumerSessions())) {
             shutdownClientGroupConsumer(clientGroupWrapper);
         }
@@ -425,19 +425,13 @@ public class ClientSessionGroupMapping {
                     log.error("say goodbye to pubSession error! {}", 
pubSession, e);
                 }
             }
-            try {
-                
Thread.sleep(eventMeshTCPServer.getEventMeshTCPConfiguration().getGracefulShutdownSleepIntervalInMills());
-            } catch (InterruptedException e) {
-                log.warn("Thread.sleep occur InterruptedException", e);
-            }
-        }
 
-        try {
-            Thread.sleep(1000);
-        } catch (InterruptedException e) {
-            log.warn("Thread.sleep occur InterruptedException", e);
+            
ThreadUtils.sleep(eventMeshTCPServer.getEventMeshTCPConfiguration().getGracefulShutdownSleepIntervalInMills(),
 TimeUnit.MILLISECONDS);
+
         }
 
+        ThreadUtils.sleep(1, TimeUnit.SECONDS);
+
         sessionTable.values().parallelStream().forEach(itr -> {
             try {
                 
EventMeshTcp2Client.serverGoodby2Client(this.eventMeshTCPServer, itr, this);
diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/rebalance/EventmeshRebalanceImpl.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/rebalance/EventmeshRebalanceImpl.java
index ad51bcb32..9041b1f52 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/rebalance/EventmeshRebalanceImpl.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/rebalance/EventmeshRebalanceImpl.java
@@ -18,6 +18,7 @@
 package org.apache.eventmesh.runtime.core.protocol.tcp.client.rebalance;
 
 import org.apache.eventmesh.api.registry.dto.EventMeshDataInfo;
+import org.apache.eventmesh.common.utils.ThreadUtils;
 import org.apache.eventmesh.runtime.boot.EventMeshTCPServer;
 import org.apache.eventmesh.runtime.constants.EventMeshConstants;
 import 
org.apache.eventmesh.runtime.core.protocol.tcp.client.EventMeshTcp2Client;
@@ -35,6 +36,7 @@ import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.concurrent.TimeUnit;
 
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -170,11 +172,7 @@ public class EventmeshRebalanceImpl implements 
EventMeshRebalanceStrategy {
             String redirectSessionAddr = 
EventMeshTcp2Client.redirectClient2NewEventMesh(eventMeshTCPServer, newProxyIp,
                     Integer.parseInt(newProxyPort), sessionList.get(i), 
eventMeshTCPServer.getClientSessionGroupMapping());
             logger.info("doRebalance,redirect sessionAddr:{}", 
redirectSessionAddr);
-            try {
-                
Thread.sleep(eventMeshTCPServer.getEventMeshTCPConfiguration().getSleepIntervalInRebalanceRedirectMills());
-            } catch (InterruptedException e) {
-                logger.warn("Thread.sleep occur InterruptedException", e);
-            }
+            
ThreadUtils.sleep(eventMeshTCPServer.getEventMeshTCPConfiguration().getSleepIntervalInRebalanceRedirectMills(),
 TimeUnit.MILLISECONDS);
         }
         logger.info("doRebalance redirect end---------------------group:{}", 
group);
     }
diff --git 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/AsyncPublishInstance.java
 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/AsyncPublishInstance.java
index 0e74b3030..a523d6d51 100644
--- 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/AsyncPublishInstance.java
+++ 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/AsyncPublishInstance.java
@@ -25,6 +25,9 @@ import org.apache.eventmesh.common.utils.IPUtils;
 import org.apache.eventmesh.common.utils.RandomStringUtils;
 import org.apache.eventmesh.common.utils.ThreadUtils;
 
+import java.util.concurrent.TimeUnit;
+
+
 import lombok.extern.slf4j.Slf4j;
 
 @Slf4j
@@ -55,9 +58,9 @@ public class AsyncPublishInstance {
                     .addProp(Constants.EVENTMESH_MESSAGE_CONST_TTL, 
String.valueOf(4 * 1000));
 
             eventMeshHttpProducer.publish(eventMeshMessage);
-            Thread.sleep(1000);
+            ThreadUtils.sleep(1, TimeUnit.SECONDS);
         }
-        Thread.sleep(30000);
+        ThreadUtils.sleep(30, TimeUnit.SECONDS);
         try (EventMeshHttpProducer ignore = eventMeshHttpProducer) {
             // ignore
         }
diff --git 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/AsyncSyncRequestInstance.java
 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/AsyncSyncRequestInstance.java
index 02c607c82..cb0543a39 100644
--- 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/AsyncSyncRequestInstance.java
+++ 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/AsyncSyncRequestInstance.java
@@ -28,6 +28,8 @@ import org.apache.eventmesh.common.utils.ThreadUtils;
 
 import org.apache.commons.lang3.StringUtils;
 
+import java.util.concurrent.TimeUnit;
+
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -79,12 +81,12 @@ public class AsyncSyncRequestInstance {
                 }
             }, 3000);
 
-            Thread.sleep(2000);
+            ThreadUtils.sleep(2, TimeUnit.SECONDS);
         } catch (Exception e) {
             logger.warn("async send msg failed", e);
         }
 
-        Thread.sleep(30000);
+        ThreadUtils.sleep(30, TimeUnit.SECONDS);
         try (final EventMeshHttpProducer ignore = eventMeshHttpProducer) {
             // close producer
         } catch (Exception e1) {
diff --git 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
index 7e4c5f01e..de12f5438 100644
--- 
a/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
+++ 
b/eventmesh-sdk-java/src/test/java/org/apache/eventmesh/client/http/demo/SyncRequestInstance.java
@@ -26,6 +26,8 @@ import org.apache.eventmesh.common.utils.ThreadUtils;
 
 import org.apache.commons.lang3.StringUtils;
 
+import java.util.concurrent.TimeUnit;
+
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -73,7 +75,7 @@ public class SyncRequestInstance {
             logger.warn("send msg failed", e);
         }
 
-        Thread.sleep(30000);
+        ThreadUtils.sleep(30, TimeUnit.SECONDS);
         try (final EventMeshHttpProducer closed = eventMeshHttpProducer) {
             // close producer
         } catch (Exception e1) {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to