Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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 {

Expand Down Expand Up @@ -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);
}
}

Expand Down Expand Up @@ -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:{}",
Expand All @@ -188,8 +198,9 @@ public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> 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;
}
Expand All @@ -203,7 +214,10 @@ public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> 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(),
Expand All @@ -222,7 +236,8 @@ public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> 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);
}
}

Expand All @@ -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<SubscriptionData> subscriptions) {
if (subscriptions == null) {
return "null";
}
List<String> 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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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";
Expand Down Expand Up @@ -433,4 +467,4 @@ public int compareTo(@NotNull ChannelId o) {
return this.channelId.compareTo(o.asLongText());
}
}
}
}
Loading