mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-28 20:09:14 +08:00
@@ -170,7 +170,7 @@ public class CommitLog implements Swappable {
|
||||
scanFileAndSetReadMode(LibC.MADV_RANDOM);
|
||||
}
|
||||
this.mappedFileQueue.checkSelf();
|
||||
log.info("load commit log " + (result ? "OK" : "Failed"));
|
||||
log.info("load commit log {}", result ? "OK" : "Failed");
|
||||
return result;
|
||||
}
|
||||
|
||||
@@ -362,14 +362,14 @@ public class CommitLog implements Swappable {
|
||||
index++;
|
||||
if (index >= mappedFiles.size()) {
|
||||
// Current branch can not happen
|
||||
log.info("recover last 3 physics file over, last mapped file " + mappedFile.getFileName());
|
||||
log.info("recover last 3 physics file over, last mapped file {}", mappedFile.getFileName());
|
||||
break;
|
||||
} else {
|
||||
mappedFile = mappedFiles.get(index);
|
||||
byteBuffer = mappedFile.sliceByteBuffer();
|
||||
processOffset = mappedFile.getFileFromOffset();
|
||||
mappedFileOffset = 0;
|
||||
log.info("recover next physics file, " + mappedFile.getFileName());
|
||||
log.info("recover next physics file, {}", mappedFile.getFileName());
|
||||
}
|
||||
}
|
||||
// Intermediate file read error
|
||||
@@ -377,7 +377,7 @@ public class CommitLog implements Swappable {
|
||||
if (size > 0) {
|
||||
log.warn("found a half message at {}, it will be truncated.", processOffset + mappedFileOffset);
|
||||
}
|
||||
log.info("recover physics file end, " + mappedFile.getFileName());
|
||||
log.info("recover physics file end, {}", mappedFile.getFileName());
|
||||
break;
|
||||
}
|
||||
}
|
||||
@@ -448,7 +448,7 @@ public class CommitLog implements Swappable {
|
||||
case BLANK_MAGIC_CODE:
|
||||
return new DispatchRequest(0, true /* success */);
|
||||
default:
|
||||
log.warn("found a illegal magic code 0x" + Integer.toHexString(magicCode));
|
||||
log.warn("found a illegal magic code 0x{}", Integer.toHexString(magicCode));
|
||||
return new DispatchRequest(-1, false /* success */);
|
||||
}
|
||||
|
||||
@@ -716,7 +716,7 @@ public class CommitLog implements Swappable {
|
||||
for (; index >= 0; index--) {
|
||||
mappedFile = mappedFiles.get(index);
|
||||
if (this.isMappedFileMatchedRecover(mappedFile, false)) {
|
||||
log.info("recover from this mapped file " + mappedFile.getFileName());
|
||||
log.info("recover from this mapped file {}", mappedFile.getFileName());
|
||||
break;
|
||||
}
|
||||
}
|
||||
@@ -761,14 +761,14 @@ public class CommitLog implements Swappable {
|
||||
if (index >= mappedFiles.size()) {
|
||||
// The current branch under normal circumstances should
|
||||
// not happen
|
||||
log.info("recover physics file over, last mapped file " + mappedFile.getFileName());
|
||||
log.info("recover physics file over, last mapped file {}", mappedFile.getFileName());
|
||||
break;
|
||||
} else {
|
||||
mappedFile = mappedFiles.get(index);
|
||||
byteBuffer = mappedFile.sliceByteBuffer();
|
||||
processOffset = mappedFile.getFileFromOffset();
|
||||
mappedFileOffset = 0;
|
||||
log.info("recover next physics file, " + mappedFile.getFileName());
|
||||
log.info("recover next physics file, {}", mappedFile.getFileName());
|
||||
}
|
||||
}
|
||||
} else {
|
||||
@@ -777,7 +777,7 @@ public class CommitLog implements Swappable {
|
||||
log.warn("found a half message at {}, it will be truncated.", processOffset + mappedFileOffset);
|
||||
}
|
||||
|
||||
log.info("recover physics file end, " + mappedFile.getFileName() + " pos=" + byteBuffer.position());
|
||||
log.info("recover physics file end, {} pos={}", mappedFile.getFileName(), byteBuffer.position());
|
||||
break;
|
||||
}
|
||||
}
|
||||
@@ -1004,7 +1004,7 @@ public class CommitLog implements Swappable {
|
||||
}
|
||||
}
|
||||
if (null == mappedFile) {
|
||||
log.error("create mapped file1 error, topic: " + msg.getTopic() + " clientAddr: " + msg.getBornHostString());
|
||||
log.error("create mapped file1 error, topic: {} clientAddr: {}", msg.getTopic(), msg.getBornHostString());
|
||||
beginTimeInLock = 0;
|
||||
return CompletableFuture.completedFuture(new PutMessageResult(PutMessageStatus.CREATE_MAPPED_FILE_FAILED, null));
|
||||
}
|
||||
@@ -1021,7 +1021,7 @@ public class CommitLog implements Swappable {
|
||||
mappedFile = this.mappedFileQueue.getLastMappedFile(0);
|
||||
if (null == mappedFile) {
|
||||
// XXX: warn and notify me
|
||||
log.error("create mapped file2 error, topic: " + msg.getTopic() + " clientAddr: " + msg.getBornHostString());
|
||||
log.error("create mapped file2 error, topic: {} clientAddr: {}", msg.getTopic(), msg.getBornHostString());
|
||||
beginTimeInLock = 0;
|
||||
return CompletableFuture.completedFuture(new PutMessageResult(PutMessageStatus.CREATE_MAPPED_FILE_FAILED, result));
|
||||
}
|
||||
@@ -1371,7 +1371,7 @@ public class CommitLog implements Swappable {
|
||||
try {
|
||||
MappedFile mappedFile = this.mappedFileQueue.getLastMappedFile(startOffset);
|
||||
if (null == mappedFile) {
|
||||
log.error("appendData getLastMappedFile error " + startOffset);
|
||||
log.error("appendData getLastMappedFile error {}", startOffset);
|
||||
return false;
|
||||
}
|
||||
|
||||
@@ -1441,7 +1441,7 @@ public class CommitLog implements Swappable {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
CommitLog.log.info(this.getServiceName() + " service started");
|
||||
CommitLog.log.info("{} service started", this.getServiceName());
|
||||
while (!this.isStopped()) {
|
||||
int interval = CommitLog.this.defaultMessageStore.getMessageStoreConfig().getCommitIntervalCommitLog();
|
||||
|
||||
@@ -1469,16 +1469,16 @@ public class CommitLog implements Swappable {
|
||||
}
|
||||
this.waitForRunning(interval);
|
||||
} catch (Throwable e) {
|
||||
CommitLog.log.error(this.getServiceName() + " service has exception. ", e);
|
||||
CommitLog.log.error("{} service has exception. ", this.getServiceName(), e);
|
||||
}
|
||||
}
|
||||
|
||||
boolean result = false;
|
||||
for (int i = 0; i < RETRY_TIMES_OVER && !result; i++) {
|
||||
result = CommitLog.this.mappedFileQueue.commit(0);
|
||||
CommitLog.log.info(this.getServiceName() + " service shutdown, retry " + (i + 1) + " times " + (result ? "OK" : "Not OK"));
|
||||
CommitLog.log.info("{} service shutdown, retry {} times {}", this.getServiceName(), i + 1, result ? "OK" : "Not OK");
|
||||
}
|
||||
CommitLog.log.info(this.getServiceName() + " service end");
|
||||
CommitLog.log.info("{} service end", this.getServiceName());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1488,7 +1488,7 @@ public class CommitLog implements Swappable {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
CommitLog.log.info(this.getServiceName() + " service started");
|
||||
CommitLog.log.info("{} service started", this.getServiceName());
|
||||
|
||||
while (!this.isStopped()) {
|
||||
boolean flushCommitLogTimed = CommitLog.this.defaultMessageStore.getMessageStoreConfig().isFlushCommitLogTimed();
|
||||
@@ -1532,7 +1532,7 @@ public class CommitLog implements Swappable {
|
||||
log.info("Flush data to disk costs {} ms", past);
|
||||
}
|
||||
} catch (Throwable e) {
|
||||
CommitLog.log.warn(this.getServiceName() + " service has exception. ", e);
|
||||
CommitLog.log.warn("{} service has exception. ", this.getServiceName(), e);
|
||||
this.printFlushProgress();
|
||||
}
|
||||
}
|
||||
@@ -1541,12 +1541,12 @@ public class CommitLog implements Swappable {
|
||||
boolean result = false;
|
||||
for (int i = 0; i < RETRY_TIMES_OVER && !result; i++) {
|
||||
result = CommitLog.this.mappedFileQueue.flush(0);
|
||||
CommitLog.log.info(this.getServiceName() + " service shutdown, retry " + (i + 1) + " times " + (result ? "OK" : "Not OK"));
|
||||
CommitLog.log.info("{} service shutdown, retry {} times {}", this.getServiceName(), i + 1, result ? "OK" : "Not OK");
|
||||
}
|
||||
|
||||
this.printFlushProgress();
|
||||
|
||||
CommitLog.log.info(this.getServiceName() + " service end");
|
||||
CommitLog.log.info("{} service end", this.getServiceName());
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -1673,14 +1673,14 @@ public class CommitLog implements Swappable {
|
||||
|
||||
@Override
|
||||
public void run() {
|
||||
CommitLog.log.info(this.getServiceName() + " service started");
|
||||
CommitLog.log.info("{} service started", this.getServiceName());
|
||||
|
||||
while (!this.isStopped()) {
|
||||
try {
|
||||
this.waitForRunning(10);
|
||||
this.doCommit();
|
||||
} catch (Exception e) {
|
||||
CommitLog.log.warn(this.getServiceName() + " service has exception. ", e);
|
||||
CommitLog.log.warn("{} service has exception. ", this.getServiceName(), e);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1695,7 +1695,7 @@ public class CommitLog implements Swappable {
|
||||
this.swapRequests();
|
||||
this.doCommit();
|
||||
|
||||
CommitLog.log.info(this.getServiceName() + " service end");
|
||||
CommitLog.log.info("{} service end", this.getServiceName());
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -1781,14 +1781,14 @@ public class CommitLog implements Swappable {
|
||||
}
|
||||
|
||||
public void run() {
|
||||
CommitLog.log.info(this.getServiceName() + " service started");
|
||||
CommitLog.log.info("{} service started", this.getServiceName());
|
||||
|
||||
while (!this.isStopped()) {
|
||||
try {
|
||||
this.waitForRunning(1);
|
||||
this.doCommit();
|
||||
} catch (Exception e) {
|
||||
CommitLog.log.warn(this.getServiceName() + " service has exception. ", e);
|
||||
CommitLog.log.warn("{} service has exception. ", this.getServiceName(), e);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1806,7 +1806,7 @@ public class CommitLog implements Swappable {
|
||||
|
||||
this.doCommit();
|
||||
|
||||
CommitLog.log.info(this.getServiceName() + " service end");
|
||||
CommitLog.log.info("{} service end", this.getServiceName());
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -1880,7 +1880,7 @@ public class CommitLog implements Swappable {
|
||||
|
||||
// Exceeds the maximum message
|
||||
if (msgLen > this.messageStoreConfig.getMaxMessageSize()) {
|
||||
log.warn("message size exceeded, msg total size: " + msgLen + ", maxMessageSize: " + this.messageStoreConfig.getMaxMessageSize());
|
||||
log.warn("message size exceeded, msg total size: {}, maxMessageSize: {}", msgLen, this.messageStoreConfig.getMaxMessageSize());
|
||||
return new AppendMessageResult(AppendMessageStatus.MESSAGE_SIZE_EXCEEDED);
|
||||
}
|
||||
|
||||
@@ -2161,7 +2161,7 @@ public class CommitLog implements Swappable {
|
||||
//flushOK=false;
|
||||
}
|
||||
if (flushStatus != PutMessageStatus.PUT_OK) {
|
||||
log.error("do groupcommit, wait for flush failed, topic: " + messageExt.getTopic() + " tags: " + messageExt.getTags() + " client address: " + messageExt.getBornHostString());
|
||||
log.error("do groupcommit, wait for flush failed, topic: {} tags: {} client address: {}", messageExt.getTopic(), messageExt.getTags(), messageExt.getBornHostString());
|
||||
putMessageResult.setPutMessageStatus(PutMessageStatus.FLUSH_DISK_TIMEOUT);
|
||||
}
|
||||
} else {
|
||||
@@ -2311,7 +2311,7 @@ public class CommitLog implements Swappable {
|
||||
long costTime = this.systemClock.now() - beginClockTimestamp;
|
||||
log.info("[{}] scanFilesInPageCache-cost {} ms.", costTime > 30 * 1000 ? "NOTIFYME" : "OK", costTime);
|
||||
} catch (Throwable e) {
|
||||
log.warn(this.getServiceName() + " service has e: {}", e);
|
||||
log.warn("{} service has e: ", this.getServiceName() , e);
|
||||
}
|
||||
}
|
||||
log.info("{} service end", this.getServiceName());
|
||||
|
||||
Reference in New Issue
Block a user