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]