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 687f8a3e1 fix issue2491
new f55787e66 Merge pull request #2492 from jonyangx/issue2491
687f8a3e1 is described below
commit 687f8a3e121fc8128b1e1b061ba87c1355383f7c
Author: jonyangx <[email protected]>
AuthorDate: Wed Dec 7 12:57:31 2022 +0800
fix issue2491
---
.../runtime/client/common/ClientConstants.java | 12 ++---
.../runtime/client/common/MessageUtils.java | 4 +-
.../runtime/client/common/RequestContext.java | 8 +--
.../eventmesh/runtime/client/common/Server.java | 10 ++--
.../eventmesh/runtime/client/common/TCPClient.java | 2 +-
.../runtime/client/impl/PubClientImpl.java | 37 +++++++++-----
.../runtime/client/impl/SubClientImpl.java | 39 +++++++++------
.../eventmesh/client/tcp/common/TcpClient.java | 58 ++++++++++++++--------
.../impl/cloudevent/CloudEventTCPPubClient.java | 6 +--
.../impl/cloudevent/CloudEventTCPSubClient.java | 2 +-
.../EventMeshMessageTCPPubClient.java | 6 +--
.../EventMeshMessageTCPSubClient.java | 2 +-
12 files changed, 112 insertions(+), 74 deletions(-)
diff --git
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/ClientConstants.java
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/ClientConstants.java
index 466e8d1e5..3f69ca086 100644
---
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/ClientConstants.java
+++
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/ClientConstants.java
@@ -20,16 +20,16 @@ package org.apache.eventmesh.runtime.client.common;
/**
* ClientConstants
*/
-public interface ClientConstants {
+public class ClientConstants {
/**
* CLIENT HEART BEAT TIME
*/
- int HEARTBEAT = 1000 * 60;
+ public static final int HEARTBEAT = 1000 * 60;
- long DEFAULT_TIMEOUT_IN_MILLISECONDS = 3000;
+ public static final long DEFAULT_TIMEOUT_IN_MILLISECONDS = 3000;
- String SYNC_TOPIC = "TEST-TOPIC-TCP-SYNC";
- String ASYNC_TOPIC = "TEST-TOPIC-TCP-ASYNC";
- String BROADCAST_TOPIC = "TEST-TOPIC-TCP-BROADCAST";
+ public static final String SYNC_TOPIC = "TEST-TOPIC-TCP-SYNC";
+ public static final String ASYNC_TOPIC = "TEST-TOPIC-TCP-ASYNC";
+ public static final String BROADCAST_TOPIC = "TEST-TOPIC-TCP-BROADCAST";
}
diff --git
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/MessageUtils.java
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/MessageUtils.java
index 72b9b357a..0d981edf9 100644
---
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/MessageUtils.java
+++
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/MessageUtils.java
@@ -144,7 +144,7 @@ public class MessageUtils {
public static UserAgent generatePubClient() {
UserAgent user = new UserAgent();
- user.setHost("127.0.0.1");
+ user.setHost("localhost");
user.setPassword(generateRandomString(8));
user.setUsername("PU4283");
user.setPath("/data/app/umg_proxy");
@@ -158,7 +158,7 @@ public class MessageUtils {
public static UserAgent generateSubServer() {
UserAgent user = new UserAgent();
- user.setHost("127.0.0.1");
+ user.setHost("localhost");
user.setPassword(generateRandomString(8));
user.setUsername("PU4283");
user.setPath("/data/app/umg_proxy");
diff --git
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/RequestContext.java
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/RequestContext.java
index 63044dbe7..e1da9ab7c 100644
---
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/RequestContext.java
+++
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/RequestContext.java
@@ -26,7 +26,7 @@ import org.slf4j.LoggerFactory;
public class RequestContext {
- private static Logger logger =
LoggerFactory.getLogger(RequestContext.class);
+ private static final Logger LOGGER =
LoggerFactory.getLogger(RequestContext.class);
private Object key;
private Package request;
@@ -78,11 +78,13 @@ public class RequestContext {
public static RequestContext context(Object key, Package request,
CountDownLatch latch) throws Exception {
RequestContext c = new RequestContext(key, request, latch);
- logger.info("_RequestContext|create|key=" + key);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info("_RequestContext|create|key=" + key);
+ }
return c;
}
- public static Object key(Package request) {
+ public static Object getHeaderSeq(Package request) {
return request.getHeader().getSeq();
}
}
diff --git
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/Server.java
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/Server.java
index 8a5d9b68b..0232f3361 100644
---
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/Server.java
+++
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/Server.java
@@ -23,7 +23,7 @@ import
org.apache.eventmesh.runtime.constants.EventMeshConstants;
public class Server {
- EventMeshServer server;
+ private transient EventMeshServer eventMeshServer;
static {
System.setProperty("proxy.home", "E:\\projects\\external-1\\proxy");
@@ -36,12 +36,12 @@ public class Server {
ConfigurationWrapper configurationWrapper =
new
ConfigurationWrapper(EventMeshConstants.EVENTMESH_CONF_HOME,
EventMeshConstants.EVENTMESH_CONF_FILE, false);
- server = new EventMeshServer(configurationWrapper);
- server.init();
- server.start();
+ eventMeshServer = new EventMeshServer(configurationWrapper);
+ eventMeshServer.init();
+ eventMeshServer.start();
}
public void shutdownAccessServer() throws Exception {
- server.shutdown();
+ eventMeshServer.shutdown();
}
}
diff --git
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/TCPClient.java
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/TCPClient.java
index 90d728a18..0a14fee35 100644
---
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/TCPClient.java
+++
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/common/TCPClient.java
@@ -120,7 +120,7 @@ public abstract class TCPClient implements Closeable {
send(msg);
return null;
} else {
- Object key = RequestContext.key(msg);
+ Object key = RequestContext.getHeaderSeq(msg);
CountDownLatch latch = new CountDownLatch(1);
RequestContext c = RequestContext.context(key, msg, latch);
if (!contexts.contains(c)) {
diff --git
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/impl/PubClientImpl.java
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/impl/PubClientImpl.java
index d5a839429..d898b64f2 100644
---
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/impl/PubClientImpl.java
+++
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/impl/PubClientImpl.java
@@ -41,7 +41,7 @@ import io.netty.channel.SimpleChannelInboundHandler;
public class PubClientImpl extends TCPClient implements PubClient {
- private final Logger publogger = LoggerFactory.getLogger(this.getClass());
+ private static final Logger LOGGER =
LoggerFactory.getLogger(PubClientImpl.class);
private final UserAgent userAgent;
@@ -61,7 +61,9 @@ public class PubClientImpl extends TCPClient implements
PubClient {
public void init() throws Exception {
open(new Handler());
hello();
- publogger.info("PubClientImpl|{}|started!", clientNo);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info("PubClientImpl|{}|started!", clientNo);
+ }
}
public void reconnect() throws Exception {
@@ -85,7 +87,10 @@ public class PubClientImpl extends TCPClient implements
PubClient {
PubClientImpl.this.reconnect();
}
Package msg = MessageUtils.heartBeat();
- publogger.debug("PubClientImpl|{}|send
heartbeat|Command={}|msg={}", clientNo, msg.getHeader().getCommand(), msg);
+ if (LOGGER.isDebugEnabled()) {
+ LOGGER.debug("PubClientImpl|{}|send
heartbeat|Command={}|msg={}",
+ clientNo, msg.getHeader().getCommand(), msg);
+ }
PubClientImpl.this.dispatcher(msg,
ClientConstants.DEFAULT_TIMEOUT_IN_MILLISECONDS);
} catch (Exception e) {
//ignore
@@ -113,7 +118,9 @@ public class PubClientImpl extends TCPClient implements
PubClient {
*/
@Override
public Package rr(Package msg, long timeout) throws Exception {
- publogger.info("PubClientImpl|{}|rr|send|Command={}|msg={}", clientNo,
Command.REQUEST_TO_SERVER, msg);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info("PubClientImpl|{}|rr|send|Command={}|msg={}",
clientNo, Command.REQUEST_TO_SERVER, msg);
+ }
return dispatcher(msg, timeout);
}
@@ -158,7 +165,9 @@ public class PubClientImpl extends TCPClient implements
PubClient {
* Send an event message, the return value is ACCESS and ACK is given
*/
public Package publish(Package msg, long timeout) throws Exception {
- publogger.info("PubClientImpl|{}|publish|send|command={}|msg={}",
clientNo, msg.getHeader().getCommand(), msg);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info("PubClientImpl|{}|publish|send|command={}|msg={}",
clientNo, msg.getHeader().getCommand(), msg);
+ }
return dispatcher(msg, timeout);
}
@@ -166,7 +175,9 @@ public class PubClientImpl extends TCPClient implements
PubClient {
* send broadcast message
*/
public Package broadcast(Package msg, long timeout) throws Exception {
- publogger.info("PubClientImpl|{}|broadcast|send|type={}|msg={}",
clientNo, msg.getHeader().getCommand(), msg);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info("PubClientImpl|{}|broadcast|send|type={}|msg={}",
clientNo, msg.getHeader().getCommand(), msg);
+ }
return dispatcher(msg, timeout);
}
@@ -179,7 +190,9 @@ public class PubClientImpl extends TCPClient implements
PubClient {
private class Handler extends SimpleChannelInboundHandler<Package> {
@Override
protected void channelRead0(ChannelHandlerContext ctx, Package msg)
throws Exception {
- publogger.info("PubClientImpl|{}|receive|type={}|msg={}",
clientNo, msg.getHeader().getCommand(), msg);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info("PubClientImpl|{}|receive|type={}|msg={}",
clientNo, msg.getHeader().getCommand(), msg);
+ }
Command cmd = msg.getHeader().getCommand();
if (callback != null) {
callback.handle(msg, ctx);
@@ -190,26 +203,26 @@ public class PubClientImpl extends TCPClient implements
PubClient {
if (cmd == Command.RESPONSE_TO_CLIENT) {
Package responseToClientAck =
MessageUtils.responseToClientAck(msg);
send(responseToClientAck);
- RequestContext context = contexts.get(RequestContext.key(msg));
+ RequestContext context =
contexts.get(RequestContext.getHeaderSeq(msg));
if (context != null) {
contexts.remove(context.getKey());
context.finish(msg);
return;
} else {
- publogger.error("msg ignored,context not found .|{}|{}",
cmd, msg);
+ LOGGER.error("msg ignored,context not found .|{}|{}", cmd,
msg);
return;
}
} else if (cmd == Command.SERVER_GOODBYE_REQUEST) {
- publogger.error("server goodby request:
---------------------------" + msg);
+ LOGGER.error("server goodby request:
---------------------------" + msg);
close();
} else {
- RequestContext context = contexts.get(RequestContext.key(msg));
+ RequestContext context =
contexts.get(RequestContext.getHeaderSeq(msg));
if (context != null) {
contexts.remove(context.getKey());
context.finish(msg);
return;
} else {
- publogger.error("msg ignored,context not found .|{}|{}",
cmd, msg);
+ LOGGER.error("msg ignored,context not found .|{}|{}", cmd,
msg);
return;
}
}
diff --git
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/impl/SubClientImpl.java
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/impl/SubClientImpl.java
index 6af8b4c0e..5a39be844 100644
---
a/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/impl/SubClientImpl.java
+++
b/eventmesh-runtime/src/test/java/org/apache/eventmesh/runtime/client/impl/SubClientImpl.java
@@ -48,15 +48,15 @@ import io.netty.channel.SimpleChannelInboundHandler;
public class SubClientImpl extends TCPClient implements SubClient {
- private static final Logger logger =
LoggerFactory.getLogger(SubClientImpl.class);
+ private static final Logger LOGGER =
LoggerFactory.getLogger(SubClientImpl.class);
- private UserAgent userAgent;
+ private transient UserAgent userAgent;
- private ReceiveMsgHook callback;
+ private transient ReceiveMsgHook callback;
- private List<SubscriptionItem> subscriptionItems = new
ArrayList<SubscriptionItem>();
+ private transient List<SubscriptionItem> subscriptionItems = new
ArrayList<SubscriptionItem>();
- private ScheduledFuture<?> task;
+ private transient ScheduledFuture<?> task;
public SubClientImpl(String accessIp, int port, UserAgent agent) {
super(accessIp, port);
@@ -70,7 +70,9 @@ public class SubClientImpl extends TCPClient implements
SubClient {
public void init() throws Exception {
open(new Handler());
hello();
- logger.info("SubClientImpl|{}|started!", clientNo);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info("SubClientImpl|{}|started!", clientNo);
+ }
}
public void reconnect() throws Exception {
@@ -103,7 +105,10 @@ public class SubClientImpl extends TCPClient implements
SubClient {
SubClientImpl.this.reconnect();
}
Package msg = MessageUtils.heartBeat();
- logger.debug("SubClientImpl|{}|send
heartbeat|Command={}|msg={}", clientNo, msg.getHeader().getCommand(), msg);
+ if (LOGGER.isDebugEnabled()) {
+ LOGGER.debug("SubClientImpl|{}|send
heartbeat|Command={}|msg={}", clientNo,
+ msg.getHeader().getCommand(), msg);
+ }
SubClientImpl.this.dispatcher(msg,
ClientConstants.DEFAULT_TIMEOUT_IN_MILLISECONDS);
} catch (Exception e) {
//ignore
@@ -122,7 +127,8 @@ public class SubClientImpl extends TCPClient implements
SubClient {
this.dispatcher(msg, ClientConstants.DEFAULT_TIMEOUT_IN_MILLISECONDS);
}
- public Package justSubscribe(String topic, SubscriptionMode
subscriptionMode, SubscriptionType subscriptionType) throws Exception {
+ public Package justSubscribe(String topic, SubscriptionMode
subscriptionMode, SubscriptionType subscriptionType)
+ throws Exception {
subscriptionItems.add(new SubscriptionItem(topic, subscriptionMode,
subscriptionType));
Package msg = MessageUtils.subscribe(topic, subscriptionMode,
subscriptionType);
return this.dispatcher(msg,
ClientConstants.DEFAULT_TIMEOUT_IN_MILLISECONDS);
@@ -197,7 +203,10 @@ public class SubClientImpl extends TCPClient implements
SubClient {
@SuppressWarnings("Duplicates")
@Override
protected void channelRead0(ChannelHandlerContext ctx, Package msg)
throws Exception {
- logger.info(SubClientImpl.class.getSimpleName() +
"|receive|command={}|msg={}", msg.getHeader().getCommand(), msg);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info(SubClientImpl.class.getSimpleName() +
"|receive|command={}|msg={}",
+ msg.getHeader().getCommand(), msg);
+ }
Command cmd = msg.getHeader().getCommand();
if (callback != null) {
callback.handle(msg, ctx);
@@ -209,33 +218,33 @@ public class SubClientImpl extends TCPClient implements
SubClient {
Package responsePKG = MessageUtils.rrResponse(msg);
send(responsePKG);
} catch (Exception e) {
- logger.info("send rr request to client ack failed");
+ LOGGER.error("send rr request to client ack failed", e);
}
} else if (cmd == Command.ASYNC_MESSAGE_TO_CLIENT) {
Package asyncAck = MessageUtils.asyncMessageAck(msg);
try {
send(asyncAck);
} catch (Exception e) {
- logger.info("send async request to client ack failed");
+ LOGGER.error("send async request to client ack failed", e);
}
} else if (cmd == Command.BROADCAST_MESSAGE_TO_CLIENT) {
Package broadcastAck = MessageUtils.broadcastMessageAck(msg);
try {
send(broadcastAck);
} catch (Exception e) {
- logger.info("send broadcast request to client ack failed");
+ LOGGER.error("send broadcast request to client ack
failed", e);
}
} else if (cmd == Command.SERVER_GOODBYE_REQUEST) {
- logger.error("server goodby request:
---------------------------" + msg);
+ LOGGER.info("server goodby request:
---------------------------" + msg);
close();
} else {
//control instruction set
- RequestContext context = contexts.get(RequestContext.key(msg));
+ RequestContext context =
contexts.get(RequestContext.getHeaderSeq(msg));
if (context != null) {
contexts.remove(context.getKey());
context.finish(msg);
} else {
- logger.error("msg ignored,context not found.|{}|{}", cmd,
msg);
+ LOGGER.error("msg ignored,context not found.|{}|{}", cmd,
msg);
}
}
}
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/TcpClient.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/TcpClient.java
index 7f37b73a8..56b25a9cd 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/TcpClient.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/common/TcpClient.java
@@ -32,6 +32,9 @@ import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
import io.netty.bootstrap.Bootstrap;
import io.netty.buffer.PooledByteBufAllocator;
import io.netty.channel.AdaptiveRecvByteBufAllocator;
@@ -51,26 +54,25 @@ import io.netty.channel.socket.nio.NioSocketChannel;
import com.google.common.base.Preconditions;
import com.google.common.util.concurrent.ThreadFactoryBuilder;
-import lombok.extern.slf4j.Slf4j;
-
-@Slf4j
public abstract class TcpClient implements Closeable {
- public final int clientNo = (new Random()).nextInt(1000);
+ private static final Logger LOGGER =
LoggerFactory.getLogger(TcpClient.class);
+
+ protected static final transient int CLIENTNO = (new
Random()).nextInt(1000);
- protected final ConcurrentHashMap<Object, RequestContext> contexts = new
ConcurrentHashMap<>();
+ protected final transient ConcurrentHashMap<Object, RequestContext>
contexts = new ConcurrentHashMap<>();
- protected final String host;
- protected final int port;
- protected final UserAgent userAgent;
+ protected final transient String host;
+ protected final transient int port;
+ protected final transient UserAgent userAgent;
- private final Bootstrap bootstrap = new Bootstrap();
+ private final transient Bootstrap bootstrap = new Bootstrap();
- private final EventLoopGroup workers = new NioEventLoopGroup();
+ private final transient EventLoopGroup workers = new NioEventLoopGroup();
- private Channel channel;
+ private transient Channel channel;
- private ScheduledFuture<?> heartTask;
+ private transient ScheduledFuture<?> heartTask;
protected static final ScheduledThreadPoolExecutor scheduler = new
ScheduledThreadPoolExecutor(
Runtime.getRuntime().availableProcessors(),
@@ -104,9 +106,10 @@ public abstract class TcpClient implements Closeable {
ChannelFuture f = bootstrap.connect(host, port).sync();
InetSocketAddress localAddress = (InetSocketAddress)
f.channel().localAddress();
channel = f.channel();
- log
- .info("connected|local={}:{}|server={}",
localAddress.getAddress().getHostAddress(), localAddress.getPort(),
- host + ":" + port);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info("connected|local={}:{}|server={}",
localAddress.getAddress().getHostAddress(),
+ localAddress.getPort(), host + ":" + port);
+ }
}
@Override
@@ -120,7 +123,10 @@ public abstract class TcpClient implements Closeable {
goodbye();
} catch (Exception e) {
Thread.currentThread().interrupt();
- log.warn("close tcp client failed.|remote address={}",
channel.remoteAddress(), e);
+
+ if (LOGGER.isWarnEnabled()) {
+ LOGGER.warn("close tcp client failed.|remote address={}",
channel.remoteAddress(), e);
+ }
}
}
@@ -134,8 +140,10 @@ public abstract class TcpClient implements Closeable {
}
Package msg = MessageUtils.heartBeat();
io(msg, EventMeshCommon.DEFAULT_TIME_OUT_MILLS);
- log.debug("heart beat start {}", msg);
- } catch (Exception ignore) {
+ if (LOGGER.isDebugEnabled()) {
+ LOGGER.debug("heart beat start {}", msg);
+ }
+ } catch (Exception e) {
// ignore
}
}, EventMeshCommon.HEARTBEAT, EventMeshCommon.HEARTBEAT,
TimeUnit.MILLISECONDS);
@@ -156,7 +164,9 @@ public abstract class TcpClient implements Closeable {
if (channel.isWritable()) {
channel.writeAndFlush(msg).addListener((ChannelFutureListener)
future -> {
if (!future.isSuccess()) {
- log.warn("send msg failed", future.cause());
+ if (LOGGER.isWarnEnabled()) {
+ LOGGER.warn("send msg failed", future.cause());
+ }
}
});
} else {
@@ -171,7 +181,9 @@ public abstract class TcpClient implements Closeable {
if (!contexts.contains(c)) {
contexts.put(key, c);
} else {
- log.info("duplicate key : {}", key);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info("duplicate key : {}", key);
+ }
}
send(msg);
if (!c.getLatch().await(timeout, TimeUnit.MILLISECONDS)) {
@@ -196,8 +208,10 @@ public abstract class TcpClient implements Closeable {
return new ChannelDuplexHandler() {
@Override
public void exceptionCaught(ChannelHandlerContext ctx, Throwable
cause) {
- log
- .info("exceptionCaught, close connection.|remote
address={}", ctx.channel().remoteAddress(), cause);
+ if (LOGGER.isInfoEnabled()) {
+ LOGGER.info("exceptionCaught, close connection.|remote
address={}",
+ ctx.channel().remoteAddress(), cause);
+ }
ctx.close();
}
};
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/cloudevent/CloudEventTCPPubClient.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/cloudevent/CloudEventTCPPubClient.java
index 9be280f56..37cf8217a 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/cloudevent/CloudEventTCPPubClient.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/cloudevent/CloudEventTCPPubClient.java
@@ -82,7 +82,7 @@ class CloudEventTCPPubClient extends TcpClient implements
EventMeshTCPPubClient<
public Package rr(CloudEvent event, long timeout) throws
EventMeshException {
try {
Package msg = MessageUtils.buildPackage(event,
Command.REQUEST_TO_SERVER);
- log.info("{}|rr|send|type={}|msg={}", clientNo, msg, msg);
+ log.info("{}|rr|send|type={}|msg={}", CLIENTNO, msg, msg);
return io(msg, timeout);
} catch (Exception ex) {
throw new EventMeshException("rr error", ex);
@@ -105,7 +105,7 @@ class CloudEventTCPPubClient extends TcpClient implements
EventMeshTCPPubClient<
try {
Package msg = MessageUtils.buildPackage(cloudEvent,
Command.ASYNC_MESSAGE_TO_SERVER);
log.info("SimplePubClientImpl cloud
event|{}|publish|send|type={}|protocol={}|msg={}",
- clientNo, msg.getHeader().getCmd(),
msg.getHeader().getProperty(Constants.PROTOCOL_TYPE), msg);
+ CLIENTNO, msg.getHeader().getCmd(),
msg.getHeader().getProperty(Constants.PROTOCOL_TYPE), msg);
return io(msg, timeout);
} catch (Exception ex) {
throw new EventMeshException("publish error", ex);
@@ -116,7 +116,7 @@ class CloudEventTCPPubClient extends TcpClient implements
EventMeshTCPPubClient<
public void broadcast(CloudEvent cloudEvent, long timeout) throws
EventMeshException {
try {
Package msg = MessageUtils.buildPackage(cloudEvent,
Command.BROADCAST_MESSAGE_TO_SERVER);
- log.info("{}|publish|send|type={}|protocol={}|msg={}", clientNo,
msg.getHeader().getCmd(),
+ log.info("{}|publish|send|type={}|protocol={}|msg={}", CLIENTNO,
msg.getHeader().getCmd(),
msg.getHeader().getProperty(Constants.PROTOCOL_TYPE), msg);
super.send(msg);
} catch (Exception ex) {
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/cloudevent/CloudEventTCPSubClient.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/cloudevent/CloudEventTCPSubClient.java
index f6eee73e0..6b81eeca4 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/cloudevent/CloudEventTCPSubClient.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/cloudevent/CloudEventTCPSubClient.java
@@ -69,7 +69,7 @@ class CloudEventTCPSubClient extends TcpClient implements
EventMeshTCPSubClient<
open(new CloudEventTCPSubHandler(contexts));
hello();
heartbeat();
- log.info("SimpleSubClientImpl|{}|started!", clientNo);
+ log.info("SimpleSubClientImpl|{}|started!", CLIENTNO);
} catch (Exception ex) {
throw new EventMeshException("Initialize
EventMeshMessageTcpSubClient error", ex);
}
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 9d7b958a0..f564cb69e 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
@@ -79,7 +79,7 @@ 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);
+ log.info("{}|rr|send|type={}|msg={}", CLIENTNO, msg, msg);
return io(msg, timeout);
} catch (Exception ex) {
throw new EventMeshException("rr error");
@@ -104,7 +104,7 @@ class EventMeshMessageTCPPubClient extends TcpClient
implements EventMeshTCPPubC
try {
Package msg = MessageUtils.buildPackage(eventMeshMessage,
Command.ASYNC_MESSAGE_TO_SERVER);
log.info("SimplePubClientImpl em
message|{}|publish|send|type={}|protocol={}|msg={}",
- clientNo, msg.getHeader().getCmd(),
+ CLIENTNO, msg.getHeader().getCmd(),
msg.getHeader().getProperty(Constants.PROTOCOL_TYPE), msg);
return io(msg, timeout);
} catch (Exception ex) {
@@ -117,7 +117,7 @@ class EventMeshMessageTCPPubClient extends TcpClient
implements EventMeshTCPPubC
try {
// todo: transform EventMeshMessage to Package
Package msg = MessageUtils.buildPackage(eventMeshMessage,
Command.BROADCAST_MESSAGE_TO_SERVER);
- log.info("{}|publish|send|type={}|protocol={}|msg={}", clientNo,
msg.getHeader().getCmd(),
+ log.info("{}|publish|send|type={}|protocol={}|msg={}", CLIENTNO,
msg.getHeader().getCmd(),
msg.getHeader().getProperty(Constants.PROTOCOL_TYPE), msg);
super.send(msg);
} catch (Exception ex) {
diff --git
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
index 008fb7e92..2d8ed03e4 100644
---
a/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
+++
b/eventmesh-sdk-java/src/main/java/org/apache/eventmesh/client/tcp/impl/eventmeshmessage/EventMeshMessageTCPSubClient.java
@@ -61,7 +61,7 @@ class EventMeshMessageTCPSubClient extends TcpClient
implements EventMeshTCPSubC
open(new EventMeshMessageTCPSubHandler(contexts));
hello();
heartbeat();
- log.info("SimpleSubClientImpl|{}|started!", clientNo);
+ log.info("SimpleSubClientImpl|{}|started!", CLIENTNO);
} catch (Exception ex) {
throw new EventMeshException("Initialize
EventMeshMessageTcpSubClient error", ex);
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]