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

mikexue 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 7402ad8ba fix issue2662
     new 31dee551d Merge pull request #2685 from jonyangx/issue2662
7402ad8ba is described below

commit 7402ad8bad795afba9e028745b1e2786b4978d18
Author: jonyangx <[email protected]>
AuthorDate: Sat Dec 24 10:51:40 2022 +0800

    fix issue2662
---
 .../protocol/http/push/HTTPMessageHandler.java     | 47 ++++++++++++----------
 1 file changed, 25 insertions(+), 22 deletions(-)

diff --git 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/HTTPMessageHandler.java
 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/HTTPMessageHandler.java
index e09fca0a6..7a4fc54d9 100644
--- 
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/HTTPMessageHandler.java
+++ 
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/http/push/HTTPMessageHandler.java
@@ -45,26 +45,29 @@ import com.google.common.collect.Sets;
 
 public class HTTPMessageHandler implements MessageHandler {
 
-    public Logger logger = LoggerFactory.getLogger(this.getClass());
+    public static final Logger LOGGER = 
LoggerFactory.getLogger(HTTPMessageHandler.class);
 
-    private EventMeshConsumer eventMeshConsumer;
+    private transient EventMeshConsumer eventMeshConsumer;
 
-    private static final ScheduledExecutorService SCHEDULER = 
ThreadPoolFactory.createSingleScheduledExecutor("eventMesh-pushMsgTimeout-");
-
-    private ThreadPoolExecutor pushExecutor;
+    private static final transient ScheduledExecutorService SCHEDULER =
+            
ThreadPoolFactory.createSingleScheduledExecutor("eventMesh-pushMsgTimeout-");
 
     private static final Integer CONSUMER_GROUP_WAITING_REQUEST_THRESHOLD = 
10000;
 
+    public static final transient Map<String, Set<AbstractHTTPPushRequest>> 
waitingRequests = Maps.newConcurrentMap();
+
+    private transient ThreadPoolExecutor pushExecutor;
+
     private void checkTimeout() {
-        waitingRequests.entrySet().stream().forEach(entry -> {
-            for (AbstractHTTPPushRequest request : entry.getValue()) {
-                request.timeout();
-                
waitingRequests.get(request.handleMsgContext.getConsumerGroup()).remove(request);
-            }
+        waitingRequests.forEach((key, value) -> {
+            value.forEach(r -> {
+                r.timeout();
+                
waitingRequests.get(r.handleMsgContext.getConsumerGroup()).remove(r);
+            });
         });
+
     }
 
-    public static final Map<String, Set<AbstractHTTPPushRequest>> 
waitingRequests = Maps.newConcurrentMap();
 
     public HTTPMessageHandler(EventMeshConsumer eventMeshConsumer) {
         this.eventMeshConsumer = eventMeshConsumer;
@@ -75,11 +78,11 @@ public class HTTPMessageHandler implements MessageHandler {
 
     @Override
     public boolean handle(final HandleMsgContext handleMsgContext) {
-        Set waitingRequests4Group = MapUtils.getObject(waitingRequests,
-            handleMsgContext.getConsumerGroup(), Sets.newConcurrentHashSet());
-        if (waitingRequests4Group.size() > 
CONSUMER_GROUP_WAITING_REQUEST_THRESHOLD) {
-            logger.warn("waitingRequests is too many, so reject, this message 
will be send back to MQ, consumerGroup:{}, threshold:{}",
-                handleMsgContext.getConsumerGroup(), 
CONSUMER_GROUP_WAITING_REQUEST_THRESHOLD);
+        if (MapUtils.getObject(waitingRequests, 
handleMsgContext.getConsumerGroup(), Sets.newConcurrentHashSet()).size()
+                > CONSUMER_GROUP_WAITING_REQUEST_THRESHOLD) {
+            LOGGER.warn("waitingRequests is too many, so reject, this message 
will be send back to MQ, "
+                            + "consumerGroup:{}, threshold:{}",
+                    handleMsgContext.getConsumerGroup(), 
CONSUMER_GROUP_WAITING_REQUEST_THRESHOLD);
             return false;
         }
 
@@ -87,12 +90,12 @@ public class HTTPMessageHandler implements MessageHandler {
             pushExecutor.submit(() -> {
                 String protocolVersion = 
Objects.requireNonNull(handleMsgContext.getEvent().getSpecVersion()).toString();
 
-                Span span = 
TraceUtils.prepareClientSpan(EventMeshUtil.getCloudEventExtensionMap(protocolVersion,
 handleMsgContext.getEvent()),
-                    
EventMeshTraceConstants.TRACE_DOWNSTREAM_EVENTMESH_CLIENT_SPAN, false);
+                Span span = 
TraceUtils.prepareClientSpan(EventMeshUtil.getCloudEventExtensionMap(protocolVersion,
+                                handleMsgContext.getEvent()),
+                        
EventMeshTraceConstants.TRACE_DOWNSTREAM_EVENTMESH_CLIENT_SPAN, false);
 
                 try {
-                    AsyncHTTPPushRequest asyncPushRequest = new 
AsyncHTTPPushRequest(handleMsgContext, waitingRequests);
-                    asyncPushRequest.tryHTTPRequest();
+                    new AsyncHTTPPushRequest(handleMsgContext, 
waitingRequests).tryHTTPRequest();
                 } finally {
                     TraceUtils.finishSpan(span, handleMsgContext.getEvent());
                 }
@@ -100,8 +103,8 @@ public class HTTPMessageHandler implements MessageHandler {
             });
             return true;
         } catch (RejectedExecutionException e) {
-            logger.warn("pushMsgThreadPoolQueue is full, so reject, current 
task size {}",
-                
handleMsgContext.getEventMeshHTTPServer().getPushMsgExecutor().getQueue().size(),
 e);
+            LOGGER.warn("pushMsgThreadPoolQueue is full, so reject, current 
task size {}",
+                    
handleMsgContext.getEventMeshHTTPServer().getPushMsgExecutor().getQueue().size(),
 e);
             return false;
         }
     }


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

Reply via email to