{"id":2958,"date":"2021-07-02T22:33:12","date_gmt":"2021-07-02T14:33:12","guid":{"rendered":"\/?p=2958"},"modified":"2021-07-06T22:33:22","modified_gmt":"2021-07-06T14:33:22","slug":"5-7-%e6%96%87%e4%bb%b6%e5%88%b7%e7%9b%98%e6%9c%ba%e5%88%b6","status":"publish","type":"post","link":"http:\/\/xinblog.ltd\/?p=2958","title":{"rendered":"5.7 \u6587\u4ef6\u5237\u76d8\u673a\u5236"},"content":{"rendered":"<p>RMQ\u7684\u5b58\u50a8\u662f\u57fa\u4e8e\u4e86JDK\u7684NIO\u7684\u5185\u5b58\u6620\u5c04\u673a\u5236,\u6d88\u606f\u5b58\u50a8\u7684\u65f6\u5019\u5148\u5c06\u6d88\u606f\u8ffd\u52a0\u5230\u5185\u5b58\u4e2d,\u7136\u540e\u6839\u636e\u914d\u7f6e\u7684\u5237\u76d8\u7b56\u7565\u5728\u4e0d\u540c\u7684\u65f6\u5019\u8fdb\u884c\u5199\u78c1\u76d8<\/p>\n<p>\u5982\u679c\u662f\u540c\u6b65\u5237\u76d8,\u76f4\u63a5commit\u4e4b\u540e,\u5229\u7528MappedByteBuffer\u7684force\u8fdb\u884c\u8f93\u76d8,\u5982\u679c\u662f\u5f02\u6b65\u5237\u76d8,\u5219\u662f\u5148\u8ffd\u52a0\u5230\u5185\u5b58\u4e2d\u4e4b\u540e\u7acb\u523b\u8fd4\u56de,\u7531\u4e00\u4e2a\u5355\u72ec\u7684\u7ebf\u7a0b\u5728\u8fdb\u884c\u78c1\u76d8\u64cd\u4f5c<\/p>\n<p>\u5728broker\u7684\u914d\u7f6e\u6587\u4ef6\u4e2d,\u652f\u6301\u914d\u7f6e\u5237\u76d8\u7684\u65b9\u5f0f,\u5206\u522b\u662fASYNC_FLUSH \u5f02\u6b65\u5237\u76d8<\/p>\n<p>SYNC_FLUSH \u540c\u6b65\u5237\u76d8,\u9ed8\u8ba4\u662f\u5f02\u6b65\u5237\u76d8<\/p>\n<p>\u6211\u4eec\u6765\u8bf4Commitlog\u7684\u5237\u76d8\u673a\u5236<\/p>\n<p>\u9996\u5148\u7684\u5237\u76d8\u5b9e\u73b0\u662fcommitLog\u7684handleDiskFlush()\u65b9\u6cd5<\/p>\n<p>\u5728putMessage\u4e2d\u89e6\u53d1\u4e86\u8fd9\u4e2a\u8c03\u7528\u5173\u7cfb<\/p>\n<p>\u9996\u5148\u662f\u6709putMessage\u4f5c\u4e3a\u5165\u53c2<\/p>\n<p>\u7136\u540e\u518dhandleDiskFlush\u4e2d,\u8fdb\u884c\u5224\u65ad\u662f\u540c\u6b65\u8fd8\u662f\u5f02\u6b65<\/p>\n<table>\n<tbody>\n<tr>\n<td>if (FlushDiskType.<em>SYNC_FLUSH <\/em>== this.defaultMessageStore.getMessageStoreConfig().getFlushDiskType()) {<\/td>\n<\/tr>\n<\/tbody>\n<\/table>\n<p>\u5728\u5176\u4e2d,\u8fdb\u884c\u83b7\u53d6\u5230\u63d0\u4ea4\u7684service<\/p>\n<p>\u521b\u5efa\u51fa\u4e00\u4e2aGroupCommitRequest,\u7528\u4e8e\u63d0\u4ea4\u7684\u6570\u636e\u7ed3\u6784,\u5728\u8fd9\u4e2a\u6570\u636e\u7ed3\u6784\u4e2d\u7684\u5c5e\u6027\u4e3b\u8981\u6709<\/p>\n<p>private final long nextOffset;<\/p>\n<p>\u504f\u79fb\u91cf<\/p>\n<p>private CompletableFuture&lt;PutMessageStatus&gt; flushOKFuture = new<\/p>\n<p>CompletableFuture&lt;&gt;();<\/p>\n<p>Completable\u7528\u4e8e\u5f02\u6b65\u6216\u8005\u540c\u6b65\u901a\u4fe1<\/p>\n<p>private final long startTimestamp = System.<em>currentTimeMillis<\/em>();<\/p>\n<p>\u5f00\u59cb\u65f6\u95f4<\/p>\n<p>private long timeoutMillis = Long.<em>MAX_VALUE<\/em>;<\/p>\n<p>\u6700\u5927\u8d85\u65f6\u65f6\u95f4<\/p>\n<p>\u5bf9\u5e94\u7684service api\u4e3a<\/p>\n<table>\n<tbody>\n<tr>\n<td>private volatile List&lt;GroupCommitRequest&gt; requestsWrite = new ArrayList&lt;GroupCommitRequest&gt;();<\/p>\n<p>private volatile List&lt;GroupCommitRequest&gt; requestsRead = new ArrayList&lt;GroupCommitRequest&gt;();<\/td>\n<\/tr>\n<\/tbody>\n<\/table>\n<p>\u4e3a\u4ec0\u4e48\u5206\u4e86\u4e24\u4e2a\u961f\u5217\u5462?<\/p>\n<p>\u8fd9\u662f\u4e00\u79cd\u5929\u7136\u7684\u907f\u514d\u8bfb\u5199\u51b2\u7a81\u7684\u65b9\u5f0f<\/p>\n<p>\u5728\u4ee3\u7801\u4e2d,\u6211\u4eec\u5c06\u4e0a\u9762\u5c01\u88c5\u597d\u7684Request\u653e\u5165\u4e86\u8fd9\u4e2aservice\u7684Write\u961f\u5217\u4e2d,\u653e\u5165\u7684\u65b9\u6cd5\u4e3a<\/p>\n<table>\n<tbody>\n<tr>\n<td>public synchronized void putRequest(final GroupCommitRequest request) {<\/p>\n<p>synchronized (this.requestsWrite) {<\/p>\n<p>this.requestsWrite.add(request);<\/p>\n<p>}<\/p>\n<p>this.wakeup();<\/p>\n<p>}<\/td>\n<\/tr>\n<\/tbody>\n<\/table>\n<p>\u540c\u6b65\u7684\u5f62\u5f0f\u63d0\u4ea4\u4efb\u52a1\u5230GroupcommitService<\/p>\n<p>\u7136\u540e\u5728\u5176\u4e2d\u8c03\u7528\u4e86wakeup\u5524\u9192\u8fd9\u4e2a\u7ebf\u7a0b<\/p>\n<table>\n<tbody>\n<tr>\n<td>public void wakeup() {<\/p>\n<p>if (hasNotified.compareAndSet(false, true)) {<\/p>\n<p>waitPoint.countDown(); \/\/ notify<\/p>\n<p>}<\/p>\n<p>}<\/td>\n<\/tr>\n<\/tbody>\n<\/table>\n<p>\u5c06waitPoint\u8bbe\u7f6e\u4e3a0,\u8fdb\u884c\u5524\u9192<\/p>\n<p>\u5728\u5524\u9192\u4e4b\u540e Service\u7684\u7236\u7c7bServiceThread\u4f1a\u8c03\u7528\u62bd\u8c61\u65b9\u6cd5 onWaitEnd,\u4ea4\u7531\u5b9e\u9645\u7684\u5b9e\u73b0\u8005<\/p>\n<p>GroupCommitService\u8fdb\u884c\u6267\u884c,\u5176\u4e2d\u8fdb\u884c\u4ea4\u6362\u8bfb\u548c\u5199\u7684\u961f\u5217<\/p>\n<table>\n<tbody>\n<tr>\n<td>private void swapRequests() {<\/p>\n<p>\/\/\u8fd9\u4e00\u6b65\u4ea4\u6362\u8bfb\u5199\u5bb9\u5668,\u611f\u89c9\u5e94\u8be5\u662f\u907f\u514d\u8bfb\u5199\u7684\u4e0d\u4e00\u81f4,\u65b9\u4fbf\u63d0\u4ea4\u7684\u65f6\u5019\u7ee7\u7eed\u653e\u5165\u63d0\u4ea4\u8bf7\u6c42<\/p>\n<p>List&lt;GroupCommitRequest&gt; tmp = this.requestsWrite;<\/p>\n<p>this.requestsWrite = this.requestsRead;<\/p>\n<p>this.requestsRead = tmp;<\/p>\n<p>}<\/td>\n<\/tr>\n<\/tbody>\n<\/table>\n<p>\u7136\u540e\u5728run\u51fd\u6570\u4e2d,\u4e00\u65e6\u5b8c\u6210\u4e86\u5524\u9192,\u5c31\u8fdb\u884c\u63d0\u4ea4\u4e86<\/p>\n<p>\u5728doCommit\u51fd\u6570\u4e2d,\u662f\u5b9e\u9645\u7684\u6267\u884c\u4ee3\u7801<\/p>\n<table>\n<tbody>\n<tr>\n<td>private void doCommit() {<\/p>\n<p>synchronized (this.requestsRead) {<\/p>\n<p>\/\/\u4ea4\u6362\u540e\u7684requestRead,\u907f\u514d\u4e86\u5e76\u53d1\u95ee\u9898<\/p>\n<p>if (!this.requestsRead.isEmpty()) {<\/p>\n<p>for (GroupCommitRequest req : this.requestsRead) {<\/p>\n<p>\/\/ There may be a message in the next file, so a maximum of<\/p>\n<p>\/\/ two times the flush<\/p>\n<p>\/\/\u67e5\u770b\u662f\u4e0d\u662f\u5df2\u7ecf\u5237\u76d8\u4e86<\/p>\n<p>boolean flushOK = CommitLog.this.mappedFileQueue.getFlushedWhere() &gt;= req.getNextOffset();<\/p>\n<p>for (int i = 0; i &lt; 2 &amp;&amp; !flushOK; i++) {<\/p>\n<p>\/\/\u5b9e\u9645\u6267\u884c\u5237\u76d8<\/p>\n<p>CommitLog.this.mappedFileQueue.flush(0);<\/p>\n<p>flushOK = CommitLog.this.mappedFileQueue.getFlushedWhere() &gt;= req.getNextOffset();<\/p>\n<p>}<\/p>\n<p>\/\/\u8fdb\u884c\u63a8\u9001<\/p>\n<p>req.wakeupCustomer(flushOK ? PutMessageStatus.<em>PUT_OK <\/em>: PutMessageStatus.<em>FLUSH_DISK_TIMEOUT<\/em>);<\/p>\n<p>}<\/p>\n<p>long storeTimestamp = CommitLog.this.mappedFileQueue.getStoreTimestamp();<\/p>\n<p>if (storeTimestamp &gt; 0) {<\/p>\n<p>\/\/\u5237\u76d8\u5b8c\u6210\u4e4b\u540e,\u66f4\u65b0\u5237\u76d8\u70b9StoreCheckpoint\u4e2d\u7684physicMsgTimeStamp<\/p>\n<p>CommitLog.this.defaultMessageStore.getStoreCheckpoint().setPhysicMsgTimestamp(storeTimestamp);<\/p>\n<p>}<\/p>\n<p>this.requestsRead.clear();<\/p>\n<p>} else {<\/p>\n<p>\/\/ Because of individual messages is set to not sync flush, it<\/p>\n<p>\/\/ will come to this process<\/p>\n<p>CommitLog.this.mappedFileQueue.flush(0);<\/p>\n<p>}<\/p>\n<p>}<\/p>\n<p>}<\/td>\n<\/tr>\n<\/tbody>\n<\/table>\n<p>\u4e0a\u8ff0\u7684\u5c31\u662f\u6574\u4e2a\u540c\u6b65\u5237\u76d8\u7684\u6d41\u7a0b\u548c\u673a\u5236<\/p>\n<p>\u63a5\u4e0b\u6765\u662fBroker\u8fdb\u884c\u5f02\u6b65\u7684\u8f93\u76d8\u673a\u5236<\/p>\n<p>\u5728\u5f02\u6b65\u7684\u5237\u76d8\u4e2d,\u5206\u4e3a\u4e86\u662f\u5426\u5f00\u542f\u5806\u5916\u5185\u5b58\u7684\u65b9\u5f0f<\/p>\n<p>\u5982\u679c\u5f00\u542f\u4e86\u5806\u5916\u5185\u5b58,\u90a3\u4e48\u5237\u76d8\u7684\u6d41\u7a0b\u4e3a<\/p>\n<p>\u73b0\u5c06\u6d88\u606f\u8ffd\u52a0\u5230ByteBuffer\u4e2d<\/p>\n<p>\u7136\u540eCommitRealTimeService\u7ebf\u7a0b\u9ed8\u8ba4\u6bcf200ms\u5c31\u8bb2ByteBuffer\u4e2d\u7684\u5185\u5bb9\u63d0\u4ea4\u5230FileChannel<\/p>\n<p>FlushReadTimeService\u6bcf500\u79d2\u8fdb\u884c\u4e00\u6b21\u5199\u5165\u78c1\u76d8,\u5b9e\u9645\u4ee3\u7801\u5982\u4e0b<\/p>\n<p>flushCommitLogService.wakeup();<\/p>\n<p>\u5728\u5185\u90e8\u5524\u9192\u4e86run\u51fd\u6570<\/p>\n<table>\n<tbody>\n<tr>\n<td>@Override<\/p>\n<p>public void run() {<\/p>\n<p>CommitLog.<em>log<\/em>.info(this.getServiceName() + &#8221; service started&#8221;);<\/p>\n<p>while (!this.isStopped()) {<\/p>\n<p>\/\/\u7b49\u5f85\u65f6\u95f4,\u4e5f\u5c31\u662f\u5524\u9192\u65f6\u95f4<\/p>\n<p>int interval = CommitLog.this.defaultMessageStore.getMessageStoreConfig().getCommitIntervalCommitLog();<\/p>\n<p>\/\/\u4e00\u6b21\u4efb\u52a1\u6700\u5c11\u63d0\u4ea4\u7684\u9875\u6570<\/p>\n<p>int commitDataLeastPages = CommitLog.this.defaultMessageStore.getMessageStoreConfig().getCommitCommitLogLeastPages();<\/p>\n<p>\/\/\u4e24\u6b21\u63d0\u4ea4\u7684\u6700\u5927\u95f4\u9694<\/p>\n<p>int commitDataThoroughInterval =<\/p>\n<p>CommitLog.this.defaultMessageStore.getMessageStoreConfig().getCommitCommitLogThoroughInterval();<\/p>\n<p>\/\/\u5ffd\u7565\u9875\u6570\u7684\u63d0\u4ea4<\/p>\n<p>long begin = System.<em>currentTimeMillis<\/em>();<\/p>\n<p>if (begin &gt;= (this.lastCommitTimestamp + commitDataThoroughInterval)) {<\/p>\n<p>this.lastCommitTimestamp = begin;<\/p>\n<p>commitDataLeastPages = 0;<\/p>\n<p>}<\/p>\n<p>try {<\/p>\n<p>\/\/\u6267\u884c\u63d0\u4ea4\u64cd\u4f5c.\u4e5f\u5c31\u662f\u5c06\u63d0\u4ea4\u6570\u636e\u63d0\u4ea4\u5230\u7269\u7406\u6587\u4ef6\u7684\u5185\u5b58\u6620\u5c04\u5185\u5b58<\/p>\n<p>boolean result = CommitLog.this.mappedFileQueue.commit(commitDataLeastPages);<\/p>\n<p>long end = System.<em>currentTimeMillis<\/em>();<\/p>\n<p>if (!result) {<\/p>\n<p>\/\/\u5982\u679c\u4e3afalse,\u5c31\u5ef6\u957f\u8f93\u76d8\u65f6\u95f4<\/p>\n<p>this.lastCommitTimestamp = end; \/\/ result = false means some data committed.<\/p>\n<p>\/\/now wake up flush thread.<\/p>\n<p>flushCommitLogService.wakeup();<\/p>\n<p>}<\/p>\n<p>if (end &#8211; begin &gt; 500) {<\/p>\n<p><em>log<\/em>.info(&#8220;Commit data to file costs {} ms&#8221;, end &#8211; begin);<\/p>\n<p>}<\/p>\n<p>this.waitForRunning(interval);<\/p>\n<p>} catch (Throwable e) {<\/p>\n<p>CommitLog.<em>log<\/em>.error(this.getServiceName() + &#8221; service has exception. &#8220;, e);<\/p>\n<p>}<\/p>\n<p>}<\/td>\n<\/tr>\n<\/tbody>\n<\/table>\n<p>\u63a5\u4e0b\u6765\u662f\u5237\u76d8\u7684\u673a\u5236<\/p>\n<table>\n<tbody>\n<tr>\n<td>public void run() {<\/p>\n<p>CommitLog.<em>log<\/em>.info(this.getServiceName() + &#8221; service started&#8221;);<\/p>\n<p>while (!this.isStopped()) {<\/p>\n<p>\/\/\u662f\u5426\u4f7f\u7528Thread.sleep\u65b9\u6cd5,\u5982\u679c\u4e3afalse \u8868\u793a\u4e3aawait\u65b9\u6cd5\u8fdb\u884c\u7b49\u5f85<\/p>\n<p>boolean flushCommitLogTimed = CommitLog.this.defaultMessageStore.getMessageStoreConfig().isFlushCommitLogTimed();<\/p>\n<p>\/\/\u5237\u76d8\u95f4\u9694<\/p>\n<p>int interval = CommitLog.this.defaultMessageStore.getMessageStoreConfig().getFlushIntervalCommitLog();<\/p>\n<p>\/\/\u4e00\u6b21\u5237\u76d8\u7684\u5305\u542b\u9875\u6570<\/p>\n<p>int flushPhysicQueueLeastPages = CommitLog.this.defaultMessageStore.getMessageStoreConfig().getFlushCommitLogLeastPages();<\/p>\n<p>\/\/\u771f\u5b9e\u5237\u76d8\u7684\u6700\u5927\u95f4\u9694<\/p>\n<p>int flushPhysicQueueThoroughInterval =<\/p>\n<p>CommitLog.this.defaultMessageStore.getMessageStoreConfig().getFlushCommitLogThoroughInterval();<\/p>\n<p>boolean printFlushProgress = false;<\/p>\n<p>\/\/ Print flush progress<\/p>\n<p>\/\/ \u8fbe\u5230\u6700\u5c0f\u5ffd\u7565\u63d0\u4ea4\u65f6\u95f4<\/p>\n<p>long currentTimeMillis = System.<em>currentTimeMillis<\/em>();<\/p>\n<p>if (currentTimeMillis &gt;= (this.lastFlushTimestamp + flushPhysicQueueThoroughInterval)) {<\/p>\n<p>this.lastFlushTimestamp = currentTimeMillis;<\/p>\n<p>flushPhysicQueueLeastPages = 0;<\/p>\n<p>printFlushProgress = (printTimes++ % 10) == 0;<\/p>\n<p>}<\/p>\n<p>try {<\/p>\n<p>if (flushCommitLogTimed) {<\/p>\n<p>Thread.<em>sleep<\/em>(interval);<\/p>\n<p>} else {<\/p>\n<p>this.waitForRunning(interval);<\/p>\n<p>}<\/p>\n<p>\/\/\u5148\u7b49\u5f85\u4e00\u5b9a\u65f6\u95f4,\u518d\u6267\u884c\u8f93\u76d8\u4efb\u52a1<\/p>\n<p>if (printFlushProgress) {<\/p>\n<p>this.printFlushProgress();<\/p>\n<p>}<\/p>\n<p>long begin = System.<em>currentTimeMillis<\/em>();<\/p>\n<p>\/\/\u8fdb\u884c\u5b9e\u9645\u7684\u5237\u65b0<\/p>\n<p>CommitLog.this.mappedFileQueue.flush(flushPhysicQueueLeastPages);<\/p>\n<p>long storeTimestamp = CommitLog.this.mappedFileQueue.getStoreTimestamp();<\/p>\n<p>if (storeTimestamp &gt; 0) {<\/p>\n<p>\/\/\u66f4\u65b0\u66f4\u65b0\u65f6\u95f4\u6233<\/p>\n<p>CommitLog.this.defaultMessageStore.getStoreCheckpoint().setPhysicMsgTimestamp(storeTimestamp);<\/p>\n<p>}<\/p>\n<p>long past = System.<em>currentTimeMillis<\/em>() &#8211; begin;<\/p>\n<p>if (past &gt; 500) {<\/p>\n<p><em>log<\/em>.info(&#8220;Flush data to disk costs {} ms&#8221;, past);<\/p>\n<p>}<\/p>\n<p>} catch (Throwable e) {<\/p>\n<p>CommitLog.<em>log<\/em>.warn(this.getServiceName() + &#8221; service has exception. &#8220;, e);<\/p>\n<p>this.printFlushProgress();<\/p>\n<p>}<\/td>\n<\/tr>\n<\/tbody>\n<\/table>\n<p>\u8fdb\u884c\u76f8\u5173\u7684\u5237\u76d8\u64cd\u4f5c<\/p>\n<p>\u8fd9\u6837\u5c31\u5b8c\u6210\u4e86\u6574\u4f53\u7684\u6587\u4ef6\u5199\u5165<\/p>\n","protected":false},"excerpt":{"rendered":"<p>RMQ\u7684\u5b58\u50a8\u662f\u57fa\u4e8e\u4e86JDK\u7684N [&hellip;]<\/p>\n","protected":false},"author":1,"featured_media":0,"comment_status":"open","ping_status":"open","sticky":false,"template":"","format":"standard","meta":{"footnotes":""},"categories":[19],"tags":[],"_links":{"self":[{"href":"http:\/\/xinblog.ltd\/index.php?rest_route=\/wp\/v2\/posts\/2958"}],"collection":[{"href":"http:\/\/xinblog.ltd\/index.php?rest_route=\/wp\/v2\/posts"}],"about":[{"href":"http:\/\/xinblog.ltd\/index.php?rest_route=\/wp\/v2\/types\/post"}],"author":[{"embeddable":true,"href":"http:\/\/xinblog.ltd\/index.php?rest_route=\/wp\/v2\/users\/1"}],"replies":[{"embeddable":true,"href":"http:\/\/xinblog.ltd\/index.php?rest_route=%2Fwp%2Fv2%2Fcomments&post=2958"}],"version-history":[{"count":0,"href":"http:\/\/xinblog.ltd\/index.php?rest_route=\/wp\/v2\/posts\/2958\/revisions"}],"wp:attachment":[{"href":"http:\/\/xinblog.ltd\/index.php?rest_route=%2Fwp%2Fv2%2Fmedia&parent=2958"}],"wp:term":[{"taxonomy":"category","embeddable":true,"href":"http:\/\/xinblog.ltd\/index.php?rest_route=%2Fwp%2Fv2%2Fcategories&post=2958"},{"taxonomy":"post_tag","embeddable":true,"href":"http:\/\/xinblog.ltd\/index.php?rest_route=%2Fwp%2Fv2%2Ftags&post=2958"}],"curies":[{"name":"wp","href":"https:\/\/api.w.org\/{rel}","templated":true}]}}