diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java index 05eb6726188..7392b0c91d4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/AbstractSystemMessageSyncer.java @@ -109,18 +109,25 @@ protected void sendSystemMessage(Object data) { Duration.ofSeconds(3).toMillis() ).whenCompleteAsync((result, throwable) -> { if (throwable != null) { - log.error("send system message failed. data: {}, topic: {}", data, getBroadcastTopicName(), throwable); + log.error("send system message failed. dataSummary: {}, topic: {}", + summarizeSystemMessageData(data), getBroadcastTopicName(), throwable); return; } if (SendStatus.SEND_OK != result.getSendStatus()) { - log.error("send system message failed. data: {}, topic: {}, sendResult:{}", data, getBroadcastTopicName(), result); + log.error("send system message failed. dataSummary: {}, topic: {}, sendResult:{}", + summarizeSystemMessageData(data), getBroadcastTopicName(), result); } }); } catch (Throwable t) { - log.error("send system message failed. data: {}, topic: {}", data, targetTopic, t); + log.error("send system message failed. dataSummary: {}, topic: {}", + summarizeSystemMessageData(data), targetTopic, t); } } + protected Object summarizeSystemMessageData(Object data) { + return data; + } + protected SendMessageRequestHeader buildSendMessageRequestHeader(Message message, String producerGroup, int queueId) { SendMessageRequestHeader requestHeader = new SendMessageRequestHeader(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncer.java index e063d79707b..4c5f7b0053b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncer.java @@ -18,6 +18,7 @@ package org.apache.rocketmq.proxy.service.sysmessage; import com.alibaba.fastjson2.JSON; +import com.google.common.base.MoreObjects; import io.netty.channel.Channel; import org.apache.rocketmq.broker.client.ClientChannelInfo; import org.apache.rocketmq.broker.client.ConsumerGroupEvent; @@ -47,6 +48,7 @@ import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; public class HeartbeatSyncer extends AbstractSystemMessageSyncer { @@ -131,16 +133,21 @@ public void onConsumerRegister(String consumerGroup, ClientChannelInfo clientCha ); data.setSubscriptionDataSet(subList); - log.debug("sync register heart beat. topic:{}, data:{}", this.getBroadcastTopicName(), data); + if (log.isDebugEnabled()) { + log.debug("sync register heart beat. topic:{}, dataSummary:{}", + this.getBroadcastTopicName(), summarizeHeartbeatData(data)); + } this.sendSystemMessage(data); } catch (Throwable t) { - log.error("heartbeat register broadcast failed. group:{}, clientChannelInfo:{}, consumeType:{}, messageModel:{}, consumeFromWhere:{}, subList:{}", - consumerGroup, clientChannelInfo, consumeType, messageModel, consumeFromWhere, subList, t); + log.error("heartbeat register broadcast failed. group:{}, clientChannelInfo:{}, consumeType:{}, messageModel:{}, consumeFromWhere:{}, subscriptionSummary:{}", + consumerGroup, clientChannelInfo, consumeType, messageModel, consumeFromWhere, + summarizeSubscriptionDataSet(subList), t); } }); } catch (Throwable t) { - log.error("heartbeat submit register broadcast failed. group:{}, clientChannelInfo:{}, consumeType:{}, messageModel:{}, consumeFromWhere:{}, subList:{}", - consumerGroup, clientChannelInfo, consumeType, messageModel, consumeFromWhere, subList, t); + log.error("heartbeat submit register broadcast failed. group:{}, clientChannelInfo:{}, consumeType:{}, messageModel:{}, consumeFromWhere:{}, subscriptionSummary:{}", + consumerGroup, clientChannelInfo, consumeType, messageModel, consumeFromWhere, + summarizeSubscriptionDataSet(subList), t); } } @@ -168,7 +175,10 @@ public void onConsumerUnRegister(String consumerGroup, ClientChannelInfo clientC remoteChannel.encode() ); - log.debug("sync unregister heart beat. topic:{}, data:{}", this.getBroadcastTopicName(), data); + if (log.isDebugEnabled()) { + log.debug("sync unregister heart beat. topic:{}, dataSummary:{}", + this.getBroadcastTopicName(), summarizeHeartbeatData(data)); + } this.sendSystemMessage(data); } catch (Throwable t) { log.error("heartbeat unregister broadcast failed. group:{}, clientChannelInfo:{}, consumeType:{}", @@ -188,8 +198,9 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo } for (MessageExt msg : msgs) { + HeartbeatSyncerData data = null; try { - HeartbeatSyncerData data = JSON.parseObject(new String(msg.getBody(), StandardCharsets.UTF_8), HeartbeatSyncerData.class); + data = JSON.parseObject(new String(msg.getBody(), StandardCharsets.UTF_8), HeartbeatSyncerData.class); if (data.getLocalProxyId().equals(localProxyId)) { continue; } @@ -203,7 +214,10 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo data.getLanguage(), data.getVersion() ); - log.debug("start process remote channel. data:{}, clientChannelInfo:{}", data, clientChannelInfo); + if (log.isDebugEnabled()) { + log.debug("start process remote channel. dataSummary:{}, clientChannelInfo:{}", + summarizeHeartbeatData(data), clientChannelInfo); + } if (data.getHeartbeatType().equals(HeartbeatType.REGISTER)) { this.consumerManager.registerConsumer( data.getGroup(), @@ -222,7 +236,8 @@ public ConsumeConcurrentlyStatus consumeMessage(List msgs, ConsumeCo ); } } catch (Throwable t) { - log.error("heartbeat consume message failed. msg:{}, data:{}", msg, new String(msg.getBody(), StandardCharsets.UTF_8), t); + log.error("heartbeat consume message failed. msg:{}, dataSummary:{}", + summarizeSystemMessage(msg), summarizeHeartbeatData(data), t); } } @@ -238,4 +253,57 @@ private String buildLocalProxyId() { private static String buildKey(String group, Channel channel) { return group + "@" + channel.id().asLongText(); } + + static String summarizeHeartbeatData(HeartbeatSyncerData data) { + if (data == null) { + return "null"; + } + return MoreObjects.toStringHelper("HeartbeatSyncerData") + .add("heartbeatType", data.getHeartbeatType()) + .add("clientId", data.getClientId()) + .add("language", data.getLanguage()) + .add("version", data.getVersion()) + .add("group", data.getGroup()) + .add("consumeType", data.getConsumeType()) + .add("messageModel", data.getMessageModel()) + .add("consumeFromWhere", data.getConsumeFromWhere()) + .add("localProxyId", data.getLocalProxyId()) + .add("channelDataPresent", data.getChannelData() != null) + .add("subscriptionSummary", summarizeSubscriptionDataSet(data.getSubscriptionDataSet())) + .toString(); + } + + static String summarizeSubscriptionDataSet(Set subscriptions) { + if (subscriptions == null) { + return "null"; + } + List topics = subscriptions.stream() + .map(SubscriptionData::getTopic) + .sorted() + .collect(Collectors.toList()); + return MoreObjects.toStringHelper("SubscriptionDataSet") + .add("count", subscriptions.size()) + .add("topics", topics) + .toString(); + } + + static String summarizeSystemMessage(MessageExt msg) { + if (msg == null) { + return "null"; + } + byte[] body = msg.getBody(); + return MoreObjects.toStringHelper("MessageExt") + .add("topic", msg.getTopic()) + .add("msgId", msg.getMsgId()) + .add("bodyBytes", body == null ? 0 : body.length) + .toString(); + } + + @Override + protected Object summarizeSystemMessageData(Object data) { + if (data instanceof HeartbeatSyncerData) { + return summarizeHeartbeatData((HeartbeatSyncerData) data); + } + return super.summarizeSystemMessageData(data); + } } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncerTest.java index 9a2c5e3437d..6e71c5be452 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/sysmessage/HeartbeatSyncerTest.java @@ -77,6 +77,7 @@ import static org.awaitility.Awaitility.await; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotSame; import static org.junit.Assert.assertSame; import static org.junit.Assert.assertTrue; @@ -217,6 +218,39 @@ public void testSyncGrpcV2Channel() throws Exception { assertSame(channelInfoList.get(0).getChannel(), syncUnRegisterChannelInfoArgumentCaptor.getValue().getChannel()); } + @Test + public void testSummarizeHeartbeatDataDoesNotExposeSubscriptionExpressions() throws Exception { + String expression = "secretTagA || secretTagB"; + HeartbeatSyncerData data = new HeartbeatSyncerData( + HeartbeatType.REGISTER, + clientId, + LanguageCode.JAVA, + 5, + "consumerGroup", + ConsumeType.CONSUME_PASSIVELY, + MessageModel.CLUSTERING, + ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET, + "proxy-0", + "raw-channel-data" + ); + data.setSubscriptionDataSet(Sets.newHashSet( + FilterAPI.buildSubscriptionData("topic-a", expression), + FilterAPI.buildSubscriptionData("topic-b", "*") + )); + + String summary = HeartbeatSyncer.summarizeHeartbeatData(data); + + assertTrue(summary.contains("heartbeatType=REGISTER")); + assertTrue(summary.contains("clientId=" + clientId)); + assertTrue(summary.contains("group=consumerGroup")); + assertTrue(summary.contains("channelDataPresent=true")); + assertTrue(summary.contains("count=2")); + assertTrue(summary.contains("topic-a")); + assertTrue(summary.contains("topic-b")); + assertFalse(summary.contains(expression)); + assertFalse(summary.contains("raw-channel-data")); + } + @Test public void testSyncRemotingChannel() throws Exception { String consumerGroup = "consumerGroup"; @@ -433,4 +467,4 @@ public int compareTo(@NotNull ChannelId o) { return this.channelId.compareTo(o.asLongText()); } } -} \ No newline at end of file +}