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 d510bddf8 fix issue2600
new 2a18ec415 Merge pull request #2617 from jonyangx/issue2600
d510bddf8 is described below
commit d510bddf8b439c4708b533cba42858e5bdb64157
Author: jonyangx <[email protected]>
AuthorDate: Sat Dec 17 09:25:40 2022 +0800
fix issue2600
---
.../tcp/client/session/send/SessionSender.java | 22 +++++++++++-----------
1 file changed, 11 insertions(+), 11 deletions(-)
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
index dab4cc834..898f08cca 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/session/send/SessionSender.java
@@ -51,16 +51,16 @@ import io.opentelemetry.api.trace.Span;
public class SessionSender {
- private final Logger messageLogger = LoggerFactory.getLogger("message");
- private final Logger logger = LoggerFactory.getLogger(SessionSender.class);
+ private static final Logger MESSAGE_LOGGER =
LoggerFactory.getLogger("message");
+ private static final Logger LOGGER =
LoggerFactory.getLogger(SessionSender.class);
- private final Session session;
+ private final transient Session session;
- public long createTime = System.currentTimeMillis();
+ public transient long createTime = System.currentTimeMillis();
- public AtomicLong upMsgs = new AtomicLong(0);
+ public transient AtomicLong upMsgs = new AtomicLong(0);
- public AtomicLong failMsgCount = new AtomicLong(0);
+ public transient AtomicLong failMsgCount = new AtomicLong(0);
private static final int TRY_PERMIT_TIME_OUT = 5;
@@ -146,12 +146,12 @@ public class SessionSender {
.getEventMesh2mqMsgNum()
.incrementAndGet();
} else {
- logger.warn("send too fast,session flow control,session:{}",
session.getClient());
+ LOGGER.warn("send too fast,session flow control,session:{}",
session.getClient());
return new EventMeshTcpSendResult(header.getSeq(),
EventMeshTcpSendStatus.SEND_TOO_FAST,
EventMeshTcpSendStatus.SEND_TOO_FAST.name());
}
} catch (Exception e) {
- logger.warn("SessionSender send failed", e);
+ LOGGER.warn("SessionSender send failed", e);
if (!(e instanceof InterruptedException)) {
upstreamBuff.release();
}
@@ -180,10 +180,10 @@ public class SessionSender {
.incrementAndGet();
Command cmd;
- if (header.getCmd().equals(Command.REQUEST_TO_SERVER)) {
+ if (Command.REQUEST_TO_SERVER == header.getCmd()) {
cmd = Command.RESPONSE_TO_CLIENT;
} else {
- messageLogger.error("invalid
message|messageHeader={}|event={}", header, event);
+ MESSAGE_LOGGER.error("invalid
message|messageHeader={}|event={}", header, event);
return;
}
event = CloudEventBuilder.from(event)
@@ -210,7 +210,7 @@ public class SessionSender {
@Override
public void onException(Throwable e) {
- messageLogger.error("exception occur while sending RR
message|user={}", session.getClient(),
+ MESSAGE_LOGGER.error("exception occur while sending RR
message|user={}", session.getClient(),
new Exception(e));
TraceUtils.finishSpanWithException(session.getContext(),
cloudEvent,
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]