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]