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 4632e4168 fix issue2652 (#2682)
4632e4168 is described below
commit 4632e4168c125dc9ebe1c48d45d5fa5d2dd5241b
Author: jonyangx <[email protected]>
AuthorDate: Wed Dec 28 09:04:57 2022 +0800
fix issue2652 (#2682)
---
.../EventMeshMessageTCPPubClient.java | 44 +++++++++++++---------
1 file changed, 26 insertions(+), 18 deletions(-)
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
index f564cb69e..743b8ad17 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPPubClient.java
@@ -44,9 +44,9 @@ import lombok.extern.slf4j.Slf4j;
@Slf4j
class EventMeshMessageTCPPubClient extends TcpClient implements
EventMeshTCPPubClient<EventMeshMessage> {
- private ReceiveMsgHook<EventMeshMessage> callback;
+ private transient ReceiveMsgHook<EventMeshMessage> callback;
- private final ConcurrentHashMap<String, AsyncRRCallback>
callbackConcurrentHashMap = new ConcurrentHashMap<>();
+ private final transient ConcurrentHashMap<String, AsyncRRCallback>
callbackConcurrentHashMap = new ConcurrentHashMap<>();
public EventMeshMessageTCPPubClient(EventMeshTCPClientConfig
eventMeshTcpClientConfig) {
super(eventMeshTcpClientConfig);
@@ -58,8 +58,8 @@ class EventMeshMessageTCPPubClient extends TcpClient
implements EventMeshTCPPubC
open(new EventMeshTCPPubHandler(contexts));
hello();
heartbeat();
- } catch (Exception ex) {
- throw new EventMeshException("Initialize
EventMeshMessageTCPPubClient error", ex);
+ } catch (Exception e) {
+ throw new EventMeshException("Initialize
EventMeshMessageTCPPubClient error", e);
}
}
@@ -69,8 +69,8 @@ class EventMeshMessageTCPPubClient extends TcpClient
implements EventMeshTCPPubC
try {
super.reconnect();
hello();
- } catch (Exception ex) {
- throw new EventMeshException("reconnect error", ex);
+ } catch (Exception e) {
+ throw new EventMeshException("reconnect error", e);
}
}
@@ -79,10 +79,12 @@ class EventMeshMessageTCPPubClient extends TcpClient
implements EventMeshTCPPubC
public Package rr(EventMeshMessage eventMeshMessage, long timeout) throws
EventMeshException {
try {
Package msg = MessageUtils.buildPackage(eventMeshMessage,
Command.REQUEST_TO_SERVER);
- log.info("{}|rr|send|type={}|msg={}", CLIENTNO, msg, msg);
+ if (log.isInfoEnabled()) {
+ log.info("{}|rr|send|type={}|msg={}", CLIENTNO, msg, msg);
+ }
return io(msg, timeout);
- } catch (Exception ex) {
- throw new EventMeshException("rr error");
+ } catch (Exception e) {
+ throw new EventMeshException("rr error", e);
}
}
@@ -93,9 +95,9 @@ class EventMeshMessageTCPPubClient extends TcpClient
implements EventMeshTCPPubC
Package msg = MessageUtils.buildPackage(eventMeshMessage,
Command.REQUEST_TO_SERVER);
super.send(msg);
this.callbackConcurrentHashMap.put((String)
RequestContext.key(msg), callback);
- } catch (Exception ex) {
+ } catch (Exception e) {
// should trigger callback?
- throw new EventMeshException("asyncRR error", ex);
+ throw new EventMeshException("asyncRR error", e);
}
}
@@ -107,8 +109,8 @@ class EventMeshMessageTCPPubClient extends TcpClient
implements EventMeshTCPPubC
CLIENTNO, msg.getHeader().getCmd(),
msg.getHeader().getProperty(Constants.PROTOCOL_TYPE), msg);
return io(msg, timeout);
- } catch (Exception ex) {
- throw new EventMeshException("publish error", ex);
+ } catch (Exception e) {
+ throw new EventMeshException("publish error", e);
}
}
@@ -120,8 +122,8 @@ class EventMeshMessageTCPPubClient extends TcpClient
implements EventMeshTCPPubC
log.info("{}|publish|send|type={}|protocol={}|msg={}", CLIENTNO,
msg.getHeader().getCmd(),
msg.getHeader().getProperty(Constants.PROTOCOL_TYPE), msg);
super.send(msg);
- } catch (Exception ex) {
- throw new EventMeshException("Broadcast message error", ex);
+ } catch (Exception e) {
+ throw new EventMeshException("Broadcast message error", e);
}
}
@@ -149,7 +151,9 @@ class EventMeshMessageTCPPubClient extends TcpClient
implements EventMeshTCPPubC
public void callback(EventMeshMessage eventMeshMessage,
ChannelHandlerContext ctx) {
if (callback != null) {
callback.handle(eventMeshMessage).ifPresent(
- responseMessage ->
ctx.writeAndFlush(MessageUtils.buildPackage(responseMessage,
Command.RESPONSE_TO_SERVER))
+ responseMessage ->
+ ctx.writeAndFlush(
+
MessageUtils.buildPackage(responseMessage, Command.RESPONSE_TO_SERVER))
);
}
}
@@ -163,8 +167,12 @@ class EventMeshMessageTCPPubClient extends TcpClient
implements EventMeshTCPPubC
public void sendResponse(Package tcpPackage) {
try {
send(tcpPackage);
- } catch (Exception exception) {
- throw new RuntimeException(exception);
+ } catch (Exception e) {
+ if (e instanceof RuntimeException) {
+ throw (RuntimeException) e;
+ } else {
+ throw new RuntimeException(e);
+ }
}
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]