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 @@ -336,14 +336,7 @@ private void initScheduledThreadPoolExecutor() {

private void initMQClientFactory() throws MQClientException {
this.mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(this.defaultLitePullConsumer, this.rpcHook);
boolean registerOK = mQClientFactory.registerConsumer(this.defaultLitePullConsumer.getConsumerGroup(), this);
if (!registerOK) {
this.serviceState = ServiceState.CREATE_JUST;

throw new MQClientException("The consumer group[" + this.defaultLitePullConsumer.getConsumerGroup()
+ "] has been created before, specify another name please." + FAQUrl.suggestTodo(FAQUrl.GROUP_NAME_DUPLICATE_URL),
null);
}
mQClientFactory.registerConsumer(this.defaultLitePullConsumer.getConsumerGroup(), this);
}

private void initRebalanceImpl() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -743,14 +743,7 @@ public synchronized void start() throws MQClientException {

this.offsetStore.load();

boolean registerOK = mQClientFactory.registerConsumer(this.defaultMQPullConsumer.getConsumerGroup(), this);
if (!registerOK) {
this.serviceState = ServiceState.CREATE_JUST;

throw new MQClientException("The consumer group[" + this.defaultMQPullConsumer.getConsumerGroup()
+ "] has been created before, specify another name please." + FAQUrl.suggestTodo(FAQUrl.GROUP_NAME_DUPLICATE_URL),
null);
}
mQClientFactory.registerConsumer(this.defaultMQPullConsumer.getConsumerGroup(), this);

mQClientFactory.start();
log.info("the consumer [{}] start OK", this.defaultMQPullConsumer.getConsumerGroup());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -985,14 +985,7 @@ public synchronized void start() throws MQClientException {
// POPTODO
this.consumeMessagePopService.start();

boolean registerOK = mQClientFactory.registerConsumer(this.defaultMQPushConsumer.getConsumerGroup(), this);
if (!registerOK) {
this.serviceState = ServiceState.CREATE_JUST;
this.consumeMessageService.shutdown(defaultMQPushConsumer.getAwaitTerminationMillisWhenShutdown());
throw new MQClientException("The consumer group[" + this.defaultMQPushConsumer.getConsumerGroup()
+ "] has been created before, specify another name please." + FAQUrl.suggestTodo(FAQUrl.GROUP_NAME_DUPLICATE_URL),
null);
}
mQClientFactory.registerConsumer(this.defaultMQPushConsumer.getConsumerGroup(), this);

mQClientFactory.start();
log.info("the consumer [{}] start OK.", this.defaultMQPushConsumer.getConsumerGroup());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1088,18 +1088,19 @@ public void shutdown() {
}
}

public synchronized boolean registerConsumer(final String group, final MQConsumerInner consumer) {
public synchronized void registerConsumer(final String group, final MQConsumerInner consumer)
throws MQClientException {
if (null == group || null == consumer) {
return false;
throw new MQClientException("consumerGroup or consumer is null", null);
}

MQConsumerInner prev = this.consumerTable.putIfAbsent(group, consumer);
if (prev != null) {
log.warn("the consumer group[" + group + "] exist already.");
return false;
throw new MQClientException(
"The consumer group[" + group + "] has been created before, specify another name please."
+ FAQUrl.suggestTodo(FAQUrl.GROUP_NAME_DUPLICATE_URL),
null);
}

return true;
}

public synchronized void unregisterConsumer(final String group) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,14 +70,10 @@ public class StoreStatsService extends ServiceThread {
private volatile LongAdder[] putMessageDistributeTime;
private volatile LongAdder[] lastPutMessageDistributeTime;
private long messageStoreBootTimestamp = System.currentTimeMillis();
private volatile long putMessageEntireTimeMax = 0;
private volatile long getMessageEntireTimeMax = 0;
// for putMessageEntireTimeMax
private ReentrantLock putLock = new ReentrantLock();
// for getMessageEntireTimeMax
private ReentrantLock getLock = new ReentrantLock();
private final AtomicLong putMessageEntireTimeMax = new AtomicLong(0);
private final AtomicLong getMessageEntireTimeMax = new AtomicLong(0);

private volatile long dispatchMaxBuffer = 0;
private final AtomicLong dispatchMaxBuffer = new AtomicLong(0);

private ReentrantLock samplingLock = new ReentrantLock();
private long lastPrintTimestamp = System.currentTimeMillis();
Expand Down Expand Up @@ -164,7 +160,7 @@ private LongAdder[] resetPutMessageDistributeTime() {
}

public long getPutMessageEntireTimeMax() {
return putMessageEntireTimeMax;
return putMessageEntireTimeMax.get();
}

public void setPutMessageEntireTimeMax(long value) {
Expand Down Expand Up @@ -213,33 +209,23 @@ else if (value < 10000) {
times[12].add(1);
}

if (value > this.putMessageEntireTimeMax) {
this.putLock.lock();
this.putMessageEntireTimeMax =
value > this.putMessageEntireTimeMax ? value : this.putMessageEntireTimeMax;
this.putLock.unlock();
}
this.putMessageEntireTimeMax.accumulateAndGet(value, Math::max);
}

public long getGetMessageEntireTimeMax() {
return getMessageEntireTimeMax;
return getMessageEntireTimeMax.get();
}

public void setGetMessageEntireTimeMax(long value) {
if (value > this.getMessageEntireTimeMax) {
this.getLock.lock();
this.getMessageEntireTimeMax =
value > this.getMessageEntireTimeMax ? value : this.getMessageEntireTimeMax;
this.getLock.unlock();
}
this.getMessageEntireTimeMax.accumulateAndGet(value, Math::max);
}

public long getDispatchMaxBuffer() {
return dispatchMaxBuffer;
return dispatchMaxBuffer.get();
}

public void setDispatchMaxBuffer(long value) {
this.dispatchMaxBuffer = value > this.dispatchMaxBuffer ? value : this.dispatchMaxBuffer;
this.dispatchMaxBuffer.accumulateAndGet(value, Math::max);
}

@Override
Expand All @@ -250,22 +236,22 @@ public String toString() {
totalTimes = 1L;
}

sb.append("\truntime: " + this.getFormatRuntime() + "\r\n");
sb.append("\tputMessageEntireTimeMax: " + this.putMessageEntireTimeMax + "\r\n");
sb.append("\tputMessageTimesTotal: " + totalTimes + "\r\n");
sb.append("\tgetPutMessageFailedTimes: " + this.getPutMessageFailedTimes() + "\r\n");
sb.append("\tputMessageSizeTotal: " + this.getPutMessageSizeTotal() + "\r\n");
sb.append("\tputMessageDistributeTime: " + this.getPutMessageDistributeTimeStringInfo(totalTimes)
+ "\r\n");
sb.append("\tputMessageAverageSize: " + (this.getPutMessageSizeTotal() / totalTimes.doubleValue())
+ "\r\n");
sb.append("\tdispatchMaxBuffer: " + this.dispatchMaxBuffer + "\r\n");
sb.append("\tgetMessageEntireTimeMax: " + this.getMessageEntireTimeMax + "\r\n");
sb.append("\tputTps: " + this.getPutTps() + "\r\n");
sb.append("\tgetFoundTps: " + this.getGetFoundTps() + "\r\n");
sb.append("\tgetMissTps: " + this.getGetMissTps() + "\r\n");
sb.append("\tgetTotalTps: " + this.getGetTotalTps() + "\r\n");
sb.append("\tgetTransferredTps: " + this.getGetTransferredTps() + "\r\n");
sb.append("\truntime: ").append(this.getFormatRuntime()).append("\r\n");
sb.append("\tputMessageEntireTimeMax: ").append(this.putMessageEntireTimeMax.get()).append("\r\n");
sb.append("\tputMessageTimesTotal: ").append(totalTimes).append("\r\n");
sb.append("\tgetPutMessageFailedTimes: ").append(this.getPutMessageFailedTimes()).append("\r\n");
sb.append("\tputMessageSizeTotal: ").append(this.getPutMessageSizeTotal()).append("\r\n");
sb.append("\tputMessageDistributeTime: ").append(this.getPutMessageDistributeTimeStringInfo(totalTimes))
.append("\r\n");
sb.append("\tputMessageAverageSize: ").append(this.getPutMessageSizeTotal() / totalTimes.doubleValue())
.append("\r\n");
sb.append("\tdispatchMaxBuffer: ").append(this.dispatchMaxBuffer.get()).append("\r\n");
sb.append("\tgetMessageEntireTimeMax: ").append(this.getMessageEntireTimeMax.get()).append("\r\n");
sb.append("\tputTps: ").append(this.getPutTps()).append("\r\n");
sb.append("\tgetFoundTps: ").append(this.getGetFoundTps()).append("\r\n");
sb.append("\tgetMissTps: ").append(this.getGetMissTps()).append("\r\n");
sb.append("\tgetTotalTps: ").append(this.getGetTotalTps()).append("\r\n");
sb.append("\tgetTransferredTps: ").append(this.getGetTransferredTps()).append("\r\n");
return sb.toString();
}

Expand Down Expand Up @@ -504,16 +490,16 @@ public HashMap<String, String> getRuntimeInfo() {

result.put("bootTimestamp", String.valueOf(this.messageStoreBootTimestamp));
result.put("runtime", this.getFormatRuntime());
result.put("putMessageEntireTimeMax", String.valueOf(this.putMessageEntireTimeMax));
result.put("putMessageEntireTimeMax", String.valueOf(this.putMessageEntireTimeMax.get()));
result.put("putMessageTimesTotal", String.valueOf(totalTimes));
result.put("putMessageFailedTimes", String.valueOf(this.putMessageFailedTimes));
result.put("putMessageSizeTotal", String.valueOf(this.getPutMessageSizeTotal()));
result.put("putMessageDistributeTime",
String.valueOf(this.getPutMessageDistributeTimeStringInfo(totalTimes)));
result.put("putMessageAverageSize",
String.valueOf(this.getPutMessageSizeTotal() / totalTimes.doubleValue()));
result.put("dispatchMaxBuffer", String.valueOf(this.dispatchMaxBuffer));
result.put("getMessageEntireTimeMax", String.valueOf(this.getMessageEntireTimeMax));
result.put("dispatchMaxBuffer", String.valueOf(this.dispatchMaxBuffer.get()));
result.put("getMessageEntireTimeMax", String.valueOf(this.getMessageEntireTimeMax.get()));
result.put("putTps", this.getPutTps());
result.put("getFoundTps", this.getGetFoundTps());
result.put("getMissTps", this.getGetMissTps());
Expand Down