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 393050c4d fix #2363 (#2478)
393050c4d is described below
commit 393050c4d8edba04e850c6309599af78ef74ee92
Author: YuanXin Hu <[email protected]>
AuthorDate: Tue Dec 6 22:08:19 2022 +0800
fix #2363 (#2478)
#2363
---
.../client/session/retry/EventMeshTcpRetryer.java | 24 +++++++++++-----------
1 file changed, 12 insertions(+), 12 deletions(-)
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/retry/EventMeshTcpRetryer.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/retry/EventMeshTcpRetryer.java
index ec2faca1a..6108d8961 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/retry/EventMeshTcpRetryer.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/retry/EventMeshTcpRetryer.java
@@ -40,11 +40,11 @@ public class EventMeshTcpRetryer {
private DelayQueue<RetryContext> retrys = new DelayQueue<RetryContext>();
private ThreadPoolExecutor pool = new ThreadPoolExecutor(3,
- 3,
- 60000,
- TimeUnit.MILLISECONDS, new ArrayBlockingQueue<Runnable>(1000),
- new EventMeshThreadFactoryImpl("eventMesh-tcp-retry", true),
- new ThreadPoolExecutor.AbortPolicy());
+ 3,
+ 60000,
+ TimeUnit.MILLISECONDS, new ArrayBlockingQueue<Runnable>(1000),
+ new EventMeshThreadFactoryImpl("eventMesh-tcp-retry", true),
+ new ThreadPoolExecutor.AbortPolicy());
private Thread dispatcher;
@@ -63,28 +63,28 @@ public class EventMeshTcpRetryer {
public void pushRetry(RetryContext retryContext) {
if (retrys.size() >=
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpMsgRetryQueueSize)
{
logger.error("pushRetry fail,retrys is too much,allow max
retryQueueSize:{}, retryTimes:{}, seq:{}, bizSeq:{}",
-
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpMsgRetryQueueSize,
retryContext.retryTimes,
- retryContext.seq,
EventMeshUtil.getMessageBizSeq(retryContext.event));
+
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpMsgRetryQueueSize,
retryContext.retryTimes,
+ retryContext.seq,
EventMeshUtil.getMessageBizSeq(retryContext.event));
return;
}
int maxRetryTimes =
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpMsgAsyncRetryTimes;
if (retryContext instanceof DownStreamMsgContext) {
DownStreamMsgContext downStreamMsgContext = (DownStreamMsgContext)
retryContext;
- maxRetryTimes =
SubscriptionType.SYNC.equals(downStreamMsgContext.subscriptionItem.getType())
- ?
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpMsgSyncRetryTimes
:
-
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpMsgAsyncRetryTimes;
+ maxRetryTimes = SubscriptionType.SYNC ==
downStreamMsgContext.subscriptionItem.getType()
+ ?
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpMsgSyncRetryTimes
:
+
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpMsgAsyncRetryTimes;
}
if (retryContext.retryTimes >= maxRetryTimes) {
logger.warn("pushRetry fail,retry over maxRetryTimes:{},
retryTimes:{}, seq:{}, bizSeq:{}", maxRetryTimes,
- retryContext.retryTimes, retryContext.seq,
EventMeshUtil.getMessageBizSeq(retryContext.event));
+ retryContext.retryTimes, retryContext.seq,
EventMeshUtil.getMessageBizSeq(retryContext.event));
return;
}
retrys.offer(retryContext);
logger.info("pushRetry success,seq:{}, retryTimes:{}, bizSeq:{}",
retryContext.seq, retryContext.retryTimes,
- EventMeshUtil.getMessageBizSeq(retryContext.event));
+ EventMeshUtil.getMessageBizSeq(retryContext.event));
}
public void init() {
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]