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 34863a214 [ISSUE #2708] Method stores return result in local before
immediately returning it [ClientSessionGroupMapping] (#2786)
34863a214 is described below
commit 34863a214ca302efdaf4454d596390501866d9ee
Author: weihubeats <[email protected]>
AuthorDate: Mon Jan 2 21:49:24 2023 +0800
[ISSUE #2708] Method stores return result in local before immediately
returning it [ClientSessionGroupMapping] (#2786)
* simplify code
* simplify code
* simplify code
---
.../client/group/ClientSessionGroupMapping.java | 159 ++++++++++-----------
1 file changed, 75 insertions(+), 84 deletions(-)
diff --git
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientSessionGroupMapping.java
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientSessionGroupMapping.java
index f5d66f7f1..93f77d20f 100644
---
a/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientSessionGroupMapping.java
+++
b/eventmesh-runtime/src/main/java/org/apache/eventmesh/runtime/core/protocol/tcp/client/group/ClientSessionGroupMapping.java
@@ -45,7 +45,6 @@ import java.util.concurrent.TimeUnit;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import io.netty.channel.ChannelFuture;
import io.netty.channel.ChannelFutureListener;
import io.netty.channel.ChannelHandlerContext;
@@ -83,8 +82,7 @@ public class ClientSessionGroupMapping {
}
public Session getSession(ChannelHandlerContext ctx) {
- Session session = getSession((InetSocketAddress)
ctx.channel().remoteAddress());
- return session;
+ return getSession((InetSocketAddress) ctx.channel().remoteAddress());
}
public Session getSession(InetSocketAddress address) {
@@ -95,7 +93,7 @@ public class ClientSessionGroupMapping {
InetSocketAddress addr = (InetSocketAddress)
ctx.channel().remoteAddress();
user.setHost(addr.getHostString());
user.setPort(addr.getPort());
- Session session = null;
+ Session session;
if (!sessionTable.containsKey(addr)) {
log.info("createSession client[{}]",
RemotingHelper.parseChannelRemoteAddr(ctx.channel()));
session = new Session(user, ctx,
eventMeshTCPServer.getEventMeshTCPConfiguration());
@@ -123,13 +121,9 @@ public class ClientSessionGroupMapping {
if (session == null) {
final String remoteAddress =
RemotingHelper.parseChannelRemoteAddr(ctx.channel());
log.info("begin to close channel to remote address[{}]",
remoteAddress);
- ctx.channel().close().addListener(new ChannelFutureListener() {
- @Override
- public void operationComplete(ChannelFuture future) throws
Exception {
- log.info("close the connection to remote address[{}]
result: {}", remoteAddress,
- future.isSuccess());
- }
- });
+ ctx.channel().close().addListener(
+ (ChannelFutureListener) future -> log.info("close the
connection to remote address[{}] result: {}", remoteAddress,
+ future.isSuccess()));
SESSION_LOGGER.info("session|close|succeed|address={}|msg={}",
addr, "no session was found");
return;
}
@@ -169,13 +163,9 @@ public class ClientSessionGroupMapping {
if (session.getContext() != null) {
log.info("begin to close channel to remote address[{}]",
remoteAddress);
- session.getContext().channel().close().addListener(new
ChannelFutureListener() {
- @Override
- public void operationComplete(ChannelFuture future) throws
Exception {
- log.info("close the connection to remote address[{}]
result: {}", remoteAddress,
- future.isSuccess());
- }
- });
+ session.getContext().channel().close().addListener(
+ (ChannelFutureListener) future -> log.info("close the
connection to remote address[{}] result: {}", remoteAddress,
+ future.isSuccess()));
}
}
}
@@ -213,7 +203,7 @@ public class ClientSessionGroupMapping {
throw new Exception("client purpose config is error");
}
- session.setClientGroupWrapper(new
WeakReference<ClientGroupWrapper>(cgw));
+ session.setClientGroupWrapper(new WeakReference<>(cgw));
}
}
@@ -243,10 +233,11 @@ public class ClientSessionGroupMapping {
}
private void startClientGroupConsumer(Session session) throws Exception {
- if (!lockMap.containsKey(session.getClient().getSubsystem())) {
- lockMap.putIfAbsent(session.getClient().getSubsystem(), new
Object());
+ String subsystem = session.getClient().getSubsystem();
+ if (!lockMap.containsKey(subsystem)) {
+ lockMap.putIfAbsent(subsystem, new Object());
}
- synchronized (lockMap.get(session.getClient().getSubsystem())) {
+ synchronized (lockMap.get(subsystem)) {
log.info("readySession session[{}]", session);
ClientGroupWrapper cgw = session.getClientGroupWrapper().get();
@@ -267,14 +258,16 @@ public class ClientSessionGroupMapping {
private void cleanClientGroupWrapperByCloseSub(Session session) throws
Exception {
cleanSubscriptionInSession(session);
-
Objects.requireNonNull(session.getClientGroupWrapper().get()).removeGroupConsumerSession(session);
+ ClientGroupWrapper clientGroupWrapper =
Objects.requireNonNull(session.getClientGroupWrapper().get());
+ clientGroupWrapper.removeGroupConsumerSession(session);
handleUnackMsgsInSession(session);
- cleanClientGroupWrapperCommon(session);
+ cleanClientGroupWrapperCommon(clientGroupWrapper);
}
private void cleanClientGroupWrapperByClosePub(Session session) throws
Exception {
-
Objects.requireNonNull(session.getClientGroupWrapper().get()).removeGroupProducerSession(session);
- cleanClientGroupWrapperCommon(session);
+ ClientGroupWrapper clientGroupWrapper =
Objects.requireNonNull(session.getClientGroupWrapper().get());
+ clientGroupWrapper.removeGroupProducerSession(session);
+ cleanClientGroupWrapperCommon(clientGroupWrapper);
}
/**
@@ -284,9 +277,10 @@ public class ClientSessionGroupMapping {
*/
private void cleanSubscriptionInSession(Session session) throws Exception {
for (SubscriptionItem item :
session.getSessionContext().subscribeTopics.values()) {
-
Objects.requireNonNull(session.getClientGroupWrapper().get()).removeSubscription(item,
session);
- if
(!Objects.requireNonNull(session.getClientGroupWrapper().get()).hasSubscription(item.getTopic()))
{
-
Objects.requireNonNull(session.getClientGroupWrapper().get()).unsubscribe(item);
+ ClientGroupWrapper clientGroupWrapper =
Objects.requireNonNull(session.getClientGroupWrapper().get());
+ clientGroupWrapper.removeSubscription(item, session);
+ if (!clientGroupWrapper.hasSubscription(item.getTopic())) {
+ clientGroupWrapper.unsubscribe(item);
}
}
}
@@ -298,7 +292,8 @@ public class ClientSessionGroupMapping {
*/
private void handleUnackMsgsInSession(Session session) {
ConcurrentHashMap<String /** seq */, DownStreamMsgContext> unAckMsg =
session.getPusher().getUnAckMsg();
- if (unAckMsg.size() > 0 &&
Objects.requireNonNull(session.getClientGroupWrapper().get()).getGroupConsumerSessions().size()
> 0) {
+ ClientGroupWrapper clientGroupWrapper =
Objects.requireNonNull(session.getClientGroupWrapper().get());
+ if (unAckMsg.size() > 0 &&
clientGroupWrapper.getGroupConsumerSessions().size() > 0) {
for (Map.Entry<String, DownStreamMsgContext> entry :
unAckMsg.entrySet()) {
DownStreamMsgContext downStreamMsgContext = entry.getValue();
if (SubscriptionMode.BROADCASTING ==
downStreamMsgContext.subscriptionItem.getMode()) {
@@ -307,10 +302,10 @@ public class ClientSessionGroupMapping {
session.getClient());
continue;
}
- Session reChooseSession =
Objects.requireNonNull(session.getClientGroupWrapper().get()).getDownstreamDispatchStrategy()
-
.select(Objects.requireNonNull(session.getClientGroupWrapper().get()).getGroup(),
+ Session reChooseSession =
clientGroupWrapper.getDownstreamDispatchStrategy()
+ .select(clientGroupWrapper.getGroup(),
downStreamMsgContext.event.getSubject(),
-
Objects.requireNonNull(session.getClientGroupWrapper().get()).groupConsumerSessions);
+ clientGroupWrapper.groupConsumerSessions);
if (reChooseSession != null) {
downStreamMsgContext.session = reChooseSession;
reChooseSession.getPusher().unAckMsg(downStreamMsgContext.seq,
downStreamMsgContext);
@@ -325,85 +320,81 @@ public class ClientSessionGroupMapping {
}
}
- private void cleanClientGroupWrapperCommon(Session session) throws
Exception {
+ private void cleanClientGroupWrapperCommon(ClientGroupWrapper
clientGroupWrapper) throws Exception {
log.info("GroupConsumerSessions size:{}",
-
Objects.requireNonNull(session.getClientGroupWrapper().get()).getGroupConsumerSessions().size());
- if
(Objects.requireNonNull(session.getClientGroupWrapper().get()).getGroupConsumerSessions().size()
== 0) {
- shutdownClientGroupConsumer(session);
+ clientGroupWrapper.getGroupConsumerSessions().size());
+ if (clientGroupWrapper.getGroupConsumerSessions().size() == 0) {
+ shutdownClientGroupConsumer(clientGroupWrapper);
}
log.info("GroupProducerSessions size:{}",
-
Objects.requireNonNull(session.getClientGroupWrapper().get()).getGroupProducerSessions().size());
- if
((Objects.requireNonNull(session.getClientGroupWrapper().get()).getGroupConsumerSessions().size()
== 0)
- &&
(Objects.requireNonNull(session.getClientGroupWrapper().get()).getGroupProducerSessions().size()
== 0)) {
- shutdownClientGroupProducer(session);
-
-
clientGroupMap.remove(Objects.requireNonNull(session.getClientGroupWrapper().get()).getGroup());
-
lockMap.remove(Objects.requireNonNull(session.getClientGroupWrapper().get()).getGroup());
- log.info("remove clientGroupWrapper group[{}]",
Objects.requireNonNull(session.getClientGroupWrapper().get()).getGroup());
+ clientGroupWrapper.getGroupProducerSessions().size());
+ if ((clientGroupWrapper.getGroupConsumerSessions().size() == 0)
+ && (clientGroupWrapper.getGroupProducerSessions().size() ==
0)) {
+ shutdownClientGroupProducer(clientGroupWrapper);
+
+ clientGroupMap.remove(clientGroupWrapper.getGroup());
+ lockMap.remove(clientGroupWrapper.getGroup());
+ log.info("remove clientGroupWrapper group[{}]",
clientGroupWrapper.getGroup());
}
}
- private void shutdownClientGroupConsumer(Session session) throws Exception
{
- if
(Objects.requireNonNull(session.getClientGroupWrapper().get()).started4Broadcast.get()
== Boolean.TRUE) {
-
Objects.requireNonNull(session.getClientGroupWrapper().get()).shutdownBroadCastConsumer();
+ private void shutdownClientGroupConsumer(ClientGroupWrapper
clientGroupWrapper) throws Exception {
+ if (clientGroupWrapper.started4Broadcast.get()) {
+ clientGroupWrapper.shutdownBroadCastConsumer();
}
- if
(Objects.requireNonNull(session.getClientGroupWrapper().get()).started4Persistent.get()
== Boolean.TRUE) {
-
Objects.requireNonNull(session.getClientGroupWrapper().get()).shutdownPersistentConsumer();
+ if (clientGroupWrapper.started4Persistent.get()) {
+ clientGroupWrapper.shutdownPersistentConsumer();
}
}
- private void shutdownClientGroupProducer(Session session) throws Exception
{
- if
(Objects.requireNonNull(session.getClientGroupWrapper().get()).producerStarted.get()
== Boolean.TRUE) {
-
Objects.requireNonNull(session.getClientGroupWrapper().get()).shutdownProducer();
+ private void shutdownClientGroupProducer(ClientGroupWrapper
clientGroupWrapper) throws Exception {
+ if (clientGroupWrapper.producerStarted.get()) {
+ clientGroupWrapper.shutdownProducer();
}
}
private void initSessionCleaner() {
eventMeshTCPServer.getScheduler().scheduleAtFixedRate(
- new Runnable() {
- @Override
- public void run() {
- for (Session tmp : sessionTable.values()) {
- if (System.currentTimeMillis() -
tmp.getLastHeartbeatTime()
- >
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpSessionExpiredInMills)
{
- try {
- log.warn("clean expired
session,client:{}", tmp.getClient());
- closeSession(tmp.getContext());
- } catch (Exception e) {
- log.error("say goodbye to session error!
{}", tmp, e);
- }
+ () -> {
+ for (Session tmp : sessionTable.values()) {
+ if (System.currentTimeMillis() - tmp.getLastHeartbeatTime()
+ >
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpSessionExpiredInMills)
{
+ try {
+ if (log.isWarnEnabled()) {
+ log.warn("clean expired session,client:{}",
tmp.getClient());
}
+ closeSession(tmp.getContext());
+ } catch (Exception e) {
+ log.error("say goodbye to session error! {}", tmp,
e);
}
}
- }, 1000,
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpSessionExpiredInMills,
+ }
+ }, 1000,
eventMeshTCPServer.getEventMeshTCPConfiguration().eventMeshTcpSessionExpiredInMills,
TimeUnit.MILLISECONDS);
}
private void initDownStreamMsgContextCleaner() {
eventMeshTCPServer.getScheduler().scheduleAtFixedRate(
- new Runnable() {
- @Override
- public void run() {
-
- //scan non-broadcast msg
- for (Session tmp : sessionTable.values()) {
- for (Map.Entry<String, DownStreamMsgContext> entry
: tmp.getPusher().getUnAckMsg().entrySet()) {
- String seqKey = entry.getKey();
- DownStreamMsgContext downStreamMsgContext =
entry.getValue();
- if (!downStreamMsgContext.isExpire()) {
- continue;
- }
- downStreamMsgContext.ackMsg();
- tmp.getPusher().getUnAckMsg().remove(seqKey);
- log.warn("remove expire downStreamMsgContext,
session:{}, topic:{}, seq:{}", tmp,
- downStreamMsgContext.event.getSubject(),
seqKey);
- }
+ () -> {
+
+ //scan non-broadcast msg
+ for (Session tmp : sessionTable.values()) {
+ for (Map.Entry<String, DownStreamMsgContext> entry :
tmp.getPusher().getUnAckMsg().entrySet()) {
+ String seqKey = entry.getKey();
+ DownStreamMsgContext downStreamMsgContext =
entry.getValue();
+ if (!downStreamMsgContext.isExpire()) {
+ continue;
}
+ downStreamMsgContext.ackMsg();
+ tmp.getPusher().getUnAckMsg().remove(seqKey);
+ log.warn("remove expire downStreamMsgContext,
session:{}, topic:{}, seq:{}", tmp,
+ downStreamMsgContext.event.getSubject(), seqKey);
}
- }, 1000, 5 * 1000, TimeUnit.MILLISECONDS);
+ }
+ }, 1000, 5 * 1000, TimeUnit.MILLISECONDS);
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]