http://git-wip-us.apache.org/repos/asf/hbase-site/blob/9118853f/devapidocs/src-html/org/apache/hadoop/hbase/mob/DefaultMobStoreFlusher.html ---------------------------------------------------------------------- diff --git a/devapidocs/src-html/org/apache/hadoop/hbase/mob/DefaultMobStoreFlusher.html b/devapidocs/src-html/org/apache/hadoop/hbase/mob/DefaultMobStoreFlusher.html index 16727c8..9cc5796 100644 --- a/devapidocs/src-html/org/apache/hadoop/hbase/mob/DefaultMobStoreFlusher.html +++ b/devapidocs/src-html/org/apache/hadoop/hbase/mob/DefaultMobStoreFlusher.html @@ -118,160 +118,157 @@ <span class="sourceLineNo">110</span> // Use a store scanner to find which rows to flush.<a name="line.110"></a> <span class="sourceLineNo">111</span> long smallestReadPoint = store.getSmallestReadPoint();<a name="line.111"></a> <span class="sourceLineNo">112</span> InternalScanner scanner = createScanner(snapshot.getScanners(), smallestReadPoint, tracker);<a name="line.112"></a> -<span class="sourceLineNo">113</span> if (scanner == null) {<a name="line.113"></a> -<span class="sourceLineNo">114</span> return result; // NULL scanner returned from coprocessor hooks means skip normal processing<a name="line.114"></a> -<span class="sourceLineNo">115</span> }<a name="line.115"></a> -<span class="sourceLineNo">116</span> StoreFileWriter writer;<a name="line.116"></a> -<span class="sourceLineNo">117</span> try {<a name="line.117"></a> -<span class="sourceLineNo">118</span> // TODO: We can fail in the below block before we complete adding this flush to<a name="line.118"></a> -<span class="sourceLineNo">119</span> // list of store files. Add cleanup of anything put on filesystem if we fail.<a name="line.119"></a> -<span class="sourceLineNo">120</span> synchronized (flushLock) {<a name="line.120"></a> -<span class="sourceLineNo">121</span> status.setStatus("Flushing " + store + ": creating writer");<a name="line.121"></a> -<span class="sourceLineNo">122</span> // Write the map out to the disk<a name="line.122"></a> -<span class="sourceLineNo">123</span> writer = store.createWriterInTmp(cellsCount, store.getColumnFamilyDescriptor().getCompressionType(),<a name="line.123"></a> -<span class="sourceLineNo">124</span> false, true, true, false);<a name="line.124"></a> -<span class="sourceLineNo">125</span> IOException e = null;<a name="line.125"></a> -<span class="sourceLineNo">126</span> try {<a name="line.126"></a> -<span class="sourceLineNo">127</span> // It's a mob store, flush the cells in a mob way. This is the difference of flushing<a name="line.127"></a> -<span class="sourceLineNo">128</span> // between a normal and a mob store.<a name="line.128"></a> -<span class="sourceLineNo">129</span> performMobFlush(snapshot, cacheFlushId, scanner, writer, status, throughputController);<a name="line.129"></a> -<span class="sourceLineNo">130</span> } catch (IOException ioe) {<a name="line.130"></a> -<span class="sourceLineNo">131</span> e = ioe;<a name="line.131"></a> -<span class="sourceLineNo">132</span> // throw the exception out<a name="line.132"></a> -<span class="sourceLineNo">133</span> throw ioe;<a name="line.133"></a> -<span class="sourceLineNo">134</span> } finally {<a name="line.134"></a> -<span class="sourceLineNo">135</span> if (e != null) {<a name="line.135"></a> -<span class="sourceLineNo">136</span> writer.close();<a name="line.136"></a> -<span class="sourceLineNo">137</span> } else {<a name="line.137"></a> -<span class="sourceLineNo">138</span> finalizeWriter(writer, cacheFlushId, status);<a name="line.138"></a> -<span class="sourceLineNo">139</span> }<a name="line.139"></a> -<span class="sourceLineNo">140</span> }<a name="line.140"></a> -<span class="sourceLineNo">141</span> }<a name="line.141"></a> -<span class="sourceLineNo">142</span> } finally {<a name="line.142"></a> -<span class="sourceLineNo">143</span> scanner.close();<a name="line.143"></a> -<span class="sourceLineNo">144</span> }<a name="line.144"></a> -<span class="sourceLineNo">145</span> LOG.info("Mob store is flushed, sequenceid=" + cacheFlushId + ", memsize="<a name="line.145"></a> -<span class="sourceLineNo">146</span> + StringUtils.TraditionalBinaryPrefix.long2String(snapshot.getDataSize(), "", 1) +<a name="line.146"></a> -<span class="sourceLineNo">147</span> ", hasBloomFilter=" + writer.hasGeneralBloom() +<a name="line.147"></a> -<span class="sourceLineNo">148</span> ", into tmp file " + writer.getPath());<a name="line.148"></a> -<span class="sourceLineNo">149</span> result.add(writer.getPath());<a name="line.149"></a> -<span class="sourceLineNo">150</span> return result;<a name="line.150"></a> -<span class="sourceLineNo">151</span> }<a name="line.151"></a> -<span class="sourceLineNo">152</span><a name="line.152"></a> -<span class="sourceLineNo">153</span> /**<a name="line.153"></a> -<span class="sourceLineNo">154</span> * Flushes the cells in the mob store.<a name="line.154"></a> -<span class="sourceLineNo">155</span> * <ol>In the mob store, the cells with PUT type might have or have no mob tags.<a name="line.155"></a> -<span class="sourceLineNo">156</span> * <li>If a cell does not have a mob tag, flushing the cell to different files depends<a name="line.156"></a> -<span class="sourceLineNo">157</span> * on the value length. If the length is larger than a threshold, it's flushed to a<a name="line.157"></a> -<span class="sourceLineNo">158</span> * mob file and the mob file is flushed to a store file in HBase. Otherwise, directly<a name="line.158"></a> -<span class="sourceLineNo">159</span> * flush the cell to a store file in HBase.</li><a name="line.159"></a> -<span class="sourceLineNo">160</span> * <li>If a cell have a mob tag, its value is a mob file name, directly flush it<a name="line.160"></a> -<span class="sourceLineNo">161</span> * to a store file in HBase.</li><a name="line.161"></a> -<span class="sourceLineNo">162</span> * </ol><a name="line.162"></a> -<span class="sourceLineNo">163</span> * @param snapshot Memstore snapshot.<a name="line.163"></a> -<span class="sourceLineNo">164</span> * @param cacheFlushId Log cache flush sequence number.<a name="line.164"></a> -<span class="sourceLineNo">165</span> * @param scanner The scanner of memstore snapshot.<a name="line.165"></a> -<span class="sourceLineNo">166</span> * @param writer The store file writer.<a name="line.166"></a> -<span class="sourceLineNo">167</span> * @param status Task that represents the flush operation and may be updated with status.<a name="line.167"></a> -<span class="sourceLineNo">168</span> * @param throughputController A controller to avoid flush too fast.<a name="line.168"></a> -<span class="sourceLineNo">169</span> * @throws IOException<a name="line.169"></a> -<span class="sourceLineNo">170</span> */<a name="line.170"></a> -<span class="sourceLineNo">171</span> protected void performMobFlush(MemStoreSnapshot snapshot, long cacheFlushId,<a name="line.171"></a> -<span class="sourceLineNo">172</span> InternalScanner scanner, StoreFileWriter writer, MonitoredTask status,<a name="line.172"></a> -<span class="sourceLineNo">173</span> ThroughputController throughputController) throws IOException {<a name="line.173"></a> -<span class="sourceLineNo">174</span> StoreFileWriter mobFileWriter = null;<a name="line.174"></a> -<span class="sourceLineNo">175</span> int compactionKVMax = conf.getInt(HConstants.COMPACTION_KV_MAX,<a name="line.175"></a> -<span class="sourceLineNo">176</span> HConstants.COMPACTION_KV_MAX_DEFAULT);<a name="line.176"></a> -<span class="sourceLineNo">177</span> long mobCount = 0;<a name="line.177"></a> -<span class="sourceLineNo">178</span> long mobSize = 0;<a name="line.178"></a> -<span class="sourceLineNo">179</span> long time = snapshot.getTimeRangeTracker().getMax();<a name="line.179"></a> -<span class="sourceLineNo">180</span> mobFileWriter = mobStore.createWriterInTmp(new Date(time), snapshot.getCellsCount(),<a name="line.180"></a> -<span class="sourceLineNo">181</span> store.getColumnFamilyDescriptor().getCompressionType(), store.getRegionInfo().getStartKey(), false);<a name="line.181"></a> -<span class="sourceLineNo">182</span> // the target path is {tableName}/.mob/{cfName}/mobFiles<a name="line.182"></a> -<span class="sourceLineNo">183</span> // the relative path is mobFiles<a name="line.183"></a> -<span class="sourceLineNo">184</span> byte[] fileName = Bytes.toBytes(mobFileWriter.getPath().getName());<a name="line.184"></a> -<span class="sourceLineNo">185</span> ScannerContext scannerContext =<a name="line.185"></a> -<span class="sourceLineNo">186</span> ScannerContext.newBuilder().setBatchLimit(compactionKVMax).build();<a name="line.186"></a> -<span class="sourceLineNo">187</span> List<Cell> cells = new ArrayList<>();<a name="line.187"></a> -<span class="sourceLineNo">188</span> boolean hasMore;<a name="line.188"></a> -<span class="sourceLineNo">189</span> String flushName = ThroughputControlUtil.getNameForThrottling(store, "flush");<a name="line.189"></a> -<span class="sourceLineNo">190</span> boolean control = throughputController != null && !store.getRegionInfo().getTable().isSystemTable();<a name="line.190"></a> -<span class="sourceLineNo">191</span> if (control) {<a name="line.191"></a> -<span class="sourceLineNo">192</span> throughputController.start(flushName);<a name="line.192"></a> -<span class="sourceLineNo">193</span> }<a name="line.193"></a> -<span class="sourceLineNo">194</span> IOException ioe = null;<a name="line.194"></a> -<span class="sourceLineNo">195</span> try {<a name="line.195"></a> -<span class="sourceLineNo">196</span> do {<a name="line.196"></a> -<span class="sourceLineNo">197</span> hasMore = scanner.next(cells, scannerContext);<a name="line.197"></a> -<span class="sourceLineNo">198</span> if (!cells.isEmpty()) {<a name="line.198"></a> -<span class="sourceLineNo">199</span> for (Cell c : cells) {<a name="line.199"></a> -<span class="sourceLineNo">200</span> // If we know that this KV is going to be included always, then let us<a name="line.200"></a> -<span class="sourceLineNo">201</span> // set its memstoreTS to 0. This will help us save space when writing to<a name="line.201"></a> -<span class="sourceLineNo">202</span> // disk.<a name="line.202"></a> -<span class="sourceLineNo">203</span> if (c.getValueLength() <= mobCellValueSizeThreshold || MobUtils.isMobReferenceCell(c)<a name="line.203"></a> -<span class="sourceLineNo">204</span> || c.getTypeByte() != KeyValue.Type.Put.getCode()) {<a name="line.204"></a> -<span class="sourceLineNo">205</span> writer.append(c);<a name="line.205"></a> -<span class="sourceLineNo">206</span> } else {<a name="line.206"></a> -<span class="sourceLineNo">207</span> // append the original keyValue in the mob file.<a name="line.207"></a> -<span class="sourceLineNo">208</span> mobFileWriter.append(c);<a name="line.208"></a> -<span class="sourceLineNo">209</span> mobSize += c.getValueLength();<a name="line.209"></a> -<span class="sourceLineNo">210</span> mobCount++;<a name="line.210"></a> -<span class="sourceLineNo">211</span><a name="line.211"></a> -<span class="sourceLineNo">212</span> // append the tags to the KeyValue.<a name="line.212"></a> -<span class="sourceLineNo">213</span> // The key is same, the value is the filename of the mob file<a name="line.213"></a> -<span class="sourceLineNo">214</span> Cell reference = MobUtils.createMobRefCell(c, fileName,<a name="line.214"></a> -<span class="sourceLineNo">215</span> this.mobStore.getRefCellTags());<a name="line.215"></a> -<span class="sourceLineNo">216</span> writer.append(reference);<a name="line.216"></a> -<span class="sourceLineNo">217</span> }<a name="line.217"></a> -<span class="sourceLineNo">218</span> int len = KeyValueUtil.length(c);<a name="line.218"></a> -<span class="sourceLineNo">219</span> if (control) {<a name="line.219"></a> -<span class="sourceLineNo">220</span> throughputController.control(flushName, len);<a name="line.220"></a> -<span class="sourceLineNo">221</span> }<a name="line.221"></a> -<span class="sourceLineNo">222</span> }<a name="line.222"></a> -<span class="sourceLineNo">223</span> cells.clear();<a name="line.223"></a> -<span class="sourceLineNo">224</span> }<a name="line.224"></a> -<span class="sourceLineNo">225</span> } while (hasMore);<a name="line.225"></a> -<span class="sourceLineNo">226</span> } catch (InterruptedException e) {<a name="line.226"></a> -<span class="sourceLineNo">227</span> ioe = new InterruptedIOException(<a name="line.227"></a> -<span class="sourceLineNo">228</span> "Interrupted while control throughput of flushing " + flushName);<a name="line.228"></a> -<span class="sourceLineNo">229</span> throw ioe;<a name="line.229"></a> -<span class="sourceLineNo">230</span> } catch (IOException e) {<a name="line.230"></a> -<span class="sourceLineNo">231</span> ioe = e;<a name="line.231"></a> -<span class="sourceLineNo">232</span> throw e;<a name="line.232"></a> -<span class="sourceLineNo">233</span> } finally {<a name="line.233"></a> -<span class="sourceLineNo">234</span> if (control) {<a name="line.234"></a> -<span class="sourceLineNo">235</span> throughputController.finish(flushName);<a name="line.235"></a> +<span class="sourceLineNo">113</span> StoreFileWriter writer;<a name="line.113"></a> +<span class="sourceLineNo">114</span> try {<a name="line.114"></a> +<span class="sourceLineNo">115</span> // TODO: We can fail in the below block before we complete adding this flush to<a name="line.115"></a> +<span class="sourceLineNo">116</span> // list of store files. Add cleanup of anything put on filesystem if we fail.<a name="line.116"></a> +<span class="sourceLineNo">117</span> synchronized (flushLock) {<a name="line.117"></a> +<span class="sourceLineNo">118</span> status.setStatus("Flushing " + store + ": creating writer");<a name="line.118"></a> +<span class="sourceLineNo">119</span> // Write the map out to the disk<a name="line.119"></a> +<span class="sourceLineNo">120</span> writer = store.createWriterInTmp(cellsCount, store.getColumnFamilyDescriptor().getCompressionType(),<a name="line.120"></a> +<span class="sourceLineNo">121</span> false, true, true, false);<a name="line.121"></a> +<span class="sourceLineNo">122</span> IOException e = null;<a name="line.122"></a> +<span class="sourceLineNo">123</span> try {<a name="line.123"></a> +<span class="sourceLineNo">124</span> // It's a mob store, flush the cells in a mob way. This is the difference of flushing<a name="line.124"></a> +<span class="sourceLineNo">125</span> // between a normal and a mob store.<a name="line.125"></a> +<span class="sourceLineNo">126</span> performMobFlush(snapshot, cacheFlushId, scanner, writer, status, throughputController);<a name="line.126"></a> +<span class="sourceLineNo">127</span> } catch (IOException ioe) {<a name="line.127"></a> +<span class="sourceLineNo">128</span> e = ioe;<a name="line.128"></a> +<span class="sourceLineNo">129</span> // throw the exception out<a name="line.129"></a> +<span class="sourceLineNo">130</span> throw ioe;<a name="line.130"></a> +<span class="sourceLineNo">131</span> } finally {<a name="line.131"></a> +<span class="sourceLineNo">132</span> if (e != null) {<a name="line.132"></a> +<span class="sourceLineNo">133</span> writer.close();<a name="line.133"></a> +<span class="sourceLineNo">134</span> } else {<a name="line.134"></a> +<span class="sourceLineNo">135</span> finalizeWriter(writer, cacheFlushId, status);<a name="line.135"></a> +<span class="sourceLineNo">136</span> }<a name="line.136"></a> +<span class="sourceLineNo">137</span> }<a name="line.137"></a> +<span class="sourceLineNo">138</span> }<a name="line.138"></a> +<span class="sourceLineNo">139</span> } finally {<a name="line.139"></a> +<span class="sourceLineNo">140</span> scanner.close();<a name="line.140"></a> +<span class="sourceLineNo">141</span> }<a name="line.141"></a> +<span class="sourceLineNo">142</span> LOG.info("Mob store is flushed, sequenceid=" + cacheFlushId + ", memsize="<a name="line.142"></a> +<span class="sourceLineNo">143</span> + StringUtils.TraditionalBinaryPrefix.long2String(snapshot.getDataSize(), "", 1) +<a name="line.143"></a> +<span class="sourceLineNo">144</span> ", hasBloomFilter=" + writer.hasGeneralBloom() +<a name="line.144"></a> +<span class="sourceLineNo">145</span> ", into tmp file " + writer.getPath());<a name="line.145"></a> +<span class="sourceLineNo">146</span> result.add(writer.getPath());<a name="line.146"></a> +<span class="sourceLineNo">147</span> return result;<a name="line.147"></a> +<span class="sourceLineNo">148</span> }<a name="line.148"></a> +<span class="sourceLineNo">149</span><a name="line.149"></a> +<span class="sourceLineNo">150</span> /**<a name="line.150"></a> +<span class="sourceLineNo">151</span> * Flushes the cells in the mob store.<a name="line.151"></a> +<span class="sourceLineNo">152</span> * <ol>In the mob store, the cells with PUT type might have or have no mob tags.<a name="line.152"></a> +<span class="sourceLineNo">153</span> * <li>If a cell does not have a mob tag, flushing the cell to different files depends<a name="line.153"></a> +<span class="sourceLineNo">154</span> * on the value length. If the length is larger than a threshold, it's flushed to a<a name="line.154"></a> +<span class="sourceLineNo">155</span> * mob file and the mob file is flushed to a store file in HBase. Otherwise, directly<a name="line.155"></a> +<span class="sourceLineNo">156</span> * flush the cell to a store file in HBase.</li><a name="line.156"></a> +<span class="sourceLineNo">157</span> * <li>If a cell have a mob tag, its value is a mob file name, directly flush it<a name="line.157"></a> +<span class="sourceLineNo">158</span> * to a store file in HBase.</li><a name="line.158"></a> +<span class="sourceLineNo">159</span> * </ol><a name="line.159"></a> +<span class="sourceLineNo">160</span> * @param snapshot Memstore snapshot.<a name="line.160"></a> +<span class="sourceLineNo">161</span> * @param cacheFlushId Log cache flush sequence number.<a name="line.161"></a> +<span class="sourceLineNo">162</span> * @param scanner The scanner of memstore snapshot.<a name="line.162"></a> +<span class="sourceLineNo">163</span> * @param writer The store file writer.<a name="line.163"></a> +<span class="sourceLineNo">164</span> * @param status Task that represents the flush operation and may be updated with status.<a name="line.164"></a> +<span class="sourceLineNo">165</span> * @param throughputController A controller to avoid flush too fast.<a name="line.165"></a> +<span class="sourceLineNo">166</span> * @throws IOException<a name="line.166"></a> +<span class="sourceLineNo">167</span> */<a name="line.167"></a> +<span class="sourceLineNo">168</span> protected void performMobFlush(MemStoreSnapshot snapshot, long cacheFlushId,<a name="line.168"></a> +<span class="sourceLineNo">169</span> InternalScanner scanner, StoreFileWriter writer, MonitoredTask status,<a name="line.169"></a> +<span class="sourceLineNo">170</span> ThroughputController throughputController) throws IOException {<a name="line.170"></a> +<span class="sourceLineNo">171</span> StoreFileWriter mobFileWriter = null;<a name="line.171"></a> +<span class="sourceLineNo">172</span> int compactionKVMax = conf.getInt(HConstants.COMPACTION_KV_MAX,<a name="line.172"></a> +<span class="sourceLineNo">173</span> HConstants.COMPACTION_KV_MAX_DEFAULT);<a name="line.173"></a> +<span class="sourceLineNo">174</span> long mobCount = 0;<a name="line.174"></a> +<span class="sourceLineNo">175</span> long mobSize = 0;<a name="line.175"></a> +<span class="sourceLineNo">176</span> long time = snapshot.getTimeRangeTracker().getMax();<a name="line.176"></a> +<span class="sourceLineNo">177</span> mobFileWriter = mobStore.createWriterInTmp(new Date(time), snapshot.getCellsCount(),<a name="line.177"></a> +<span class="sourceLineNo">178</span> store.getColumnFamilyDescriptor().getCompressionType(), store.getRegionInfo().getStartKey(), false);<a name="line.178"></a> +<span class="sourceLineNo">179</span> // the target path is {tableName}/.mob/{cfName}/mobFiles<a name="line.179"></a> +<span class="sourceLineNo">180</span> // the relative path is mobFiles<a name="line.180"></a> +<span class="sourceLineNo">181</span> byte[] fileName = Bytes.toBytes(mobFileWriter.getPath().getName());<a name="line.181"></a> +<span class="sourceLineNo">182</span> ScannerContext scannerContext =<a name="line.182"></a> +<span class="sourceLineNo">183</span> ScannerContext.newBuilder().setBatchLimit(compactionKVMax).build();<a name="line.183"></a> +<span class="sourceLineNo">184</span> List<Cell> cells = new ArrayList<>();<a name="line.184"></a> +<span class="sourceLineNo">185</span> boolean hasMore;<a name="line.185"></a> +<span class="sourceLineNo">186</span> String flushName = ThroughputControlUtil.getNameForThrottling(store, "flush");<a name="line.186"></a> +<span class="sourceLineNo">187</span> boolean control = throughputController != null && !store.getRegionInfo().getTable().isSystemTable();<a name="line.187"></a> +<span class="sourceLineNo">188</span> if (control) {<a name="line.188"></a> +<span class="sourceLineNo">189</span> throughputController.start(flushName);<a name="line.189"></a> +<span class="sourceLineNo">190</span> }<a name="line.190"></a> +<span class="sourceLineNo">191</span> IOException ioe = null;<a name="line.191"></a> +<span class="sourceLineNo">192</span> try {<a name="line.192"></a> +<span class="sourceLineNo">193</span> do {<a name="line.193"></a> +<span class="sourceLineNo">194</span> hasMore = scanner.next(cells, scannerContext);<a name="line.194"></a> +<span class="sourceLineNo">195</span> if (!cells.isEmpty()) {<a name="line.195"></a> +<span class="sourceLineNo">196</span> for (Cell c : cells) {<a name="line.196"></a> +<span class="sourceLineNo">197</span> // If we know that this KV is going to be included always, then let us<a name="line.197"></a> +<span class="sourceLineNo">198</span> // set its memstoreTS to 0. This will help us save space when writing to<a name="line.198"></a> +<span class="sourceLineNo">199</span> // disk.<a name="line.199"></a> +<span class="sourceLineNo">200</span> if (c.getValueLength() <= mobCellValueSizeThreshold || MobUtils.isMobReferenceCell(c)<a name="line.200"></a> +<span class="sourceLineNo">201</span> || c.getTypeByte() != KeyValue.Type.Put.getCode()) {<a name="line.201"></a> +<span class="sourceLineNo">202</span> writer.append(c);<a name="line.202"></a> +<span class="sourceLineNo">203</span> } else {<a name="line.203"></a> +<span class="sourceLineNo">204</span> // append the original keyValue in the mob file.<a name="line.204"></a> +<span class="sourceLineNo">205</span> mobFileWriter.append(c);<a name="line.205"></a> +<span class="sourceLineNo">206</span> mobSize += c.getValueLength();<a name="line.206"></a> +<span class="sourceLineNo">207</span> mobCount++;<a name="line.207"></a> +<span class="sourceLineNo">208</span><a name="line.208"></a> +<span class="sourceLineNo">209</span> // append the tags to the KeyValue.<a name="line.209"></a> +<span class="sourceLineNo">210</span> // The key is same, the value is the filename of the mob file<a name="line.210"></a> +<span class="sourceLineNo">211</span> Cell reference = MobUtils.createMobRefCell(c, fileName,<a name="line.211"></a> +<span class="sourceLineNo">212</span> this.mobStore.getRefCellTags());<a name="line.212"></a> +<span class="sourceLineNo">213</span> writer.append(reference);<a name="line.213"></a> +<span class="sourceLineNo">214</span> }<a name="line.214"></a> +<span class="sourceLineNo">215</span> int len = KeyValueUtil.length(c);<a name="line.215"></a> +<span class="sourceLineNo">216</span> if (control) {<a name="line.216"></a> +<span class="sourceLineNo">217</span> throughputController.control(flushName, len);<a name="line.217"></a> +<span class="sourceLineNo">218</span> }<a name="line.218"></a> +<span class="sourceLineNo">219</span> }<a name="line.219"></a> +<span class="sourceLineNo">220</span> cells.clear();<a name="line.220"></a> +<span class="sourceLineNo">221</span> }<a name="line.221"></a> +<span class="sourceLineNo">222</span> } while (hasMore);<a name="line.222"></a> +<span class="sourceLineNo">223</span> } catch (InterruptedException e) {<a name="line.223"></a> +<span class="sourceLineNo">224</span> ioe = new InterruptedIOException(<a name="line.224"></a> +<span class="sourceLineNo">225</span> "Interrupted while control throughput of flushing " + flushName);<a name="line.225"></a> +<span class="sourceLineNo">226</span> throw ioe;<a name="line.226"></a> +<span class="sourceLineNo">227</span> } catch (IOException e) {<a name="line.227"></a> +<span class="sourceLineNo">228</span> ioe = e;<a name="line.228"></a> +<span class="sourceLineNo">229</span> throw e;<a name="line.229"></a> +<span class="sourceLineNo">230</span> } finally {<a name="line.230"></a> +<span class="sourceLineNo">231</span> if (control) {<a name="line.231"></a> +<span class="sourceLineNo">232</span> throughputController.finish(flushName);<a name="line.232"></a> +<span class="sourceLineNo">233</span> }<a name="line.233"></a> +<span class="sourceLineNo">234</span> if (ioe != null) {<a name="line.234"></a> +<span class="sourceLineNo">235</span> mobFileWriter.close();<a name="line.235"></a> <span class="sourceLineNo">236</span> }<a name="line.236"></a> -<span class="sourceLineNo">237</span> if (ioe != null) {<a name="line.237"></a> -<span class="sourceLineNo">238</span> mobFileWriter.close();<a name="line.238"></a> -<span class="sourceLineNo">239</span> }<a name="line.239"></a> -<span class="sourceLineNo">240</span> }<a name="line.240"></a> -<span class="sourceLineNo">241</span><a name="line.241"></a> -<span class="sourceLineNo">242</span> if (mobCount > 0) {<a name="line.242"></a> -<span class="sourceLineNo">243</span> // commit the mob file from temp folder to target folder.<a name="line.243"></a> -<span class="sourceLineNo">244</span> // If the mob file is committed successfully but the store file is not,<a name="line.244"></a> -<span class="sourceLineNo">245</span> // the committed mob file will be handled by the sweep tool as an unused<a name="line.245"></a> -<span class="sourceLineNo">246</span> // file.<a name="line.246"></a> -<span class="sourceLineNo">247</span> status.setStatus("Flushing mob file " + store + ": appending metadata");<a name="line.247"></a> -<span class="sourceLineNo">248</span> mobFileWriter.appendMetadata(cacheFlushId, false, mobCount);<a name="line.248"></a> -<span class="sourceLineNo">249</span> status.setStatus("Flushing mob file " + store + ": closing flushed file");<a name="line.249"></a> -<span class="sourceLineNo">250</span> mobFileWriter.close();<a name="line.250"></a> -<span class="sourceLineNo">251</span> mobStore.commitFile(mobFileWriter.getPath(), targetPath);<a name="line.251"></a> -<span class="sourceLineNo">252</span> mobStore.updateMobFlushCount();<a name="line.252"></a> -<span class="sourceLineNo">253</span> mobStore.updateMobFlushedCellsCount(mobCount);<a name="line.253"></a> -<span class="sourceLineNo">254</span> mobStore.updateMobFlushedCellsSize(mobSize);<a name="line.254"></a> -<span class="sourceLineNo">255</span> } else {<a name="line.255"></a> -<span class="sourceLineNo">256</span> try {<a name="line.256"></a> -<span class="sourceLineNo">257</span> status.setStatus("Flushing mob file " + store + ": no mob cells, closing flushed file");<a name="line.257"></a> -<span class="sourceLineNo">258</span> mobFileWriter.close();<a name="line.258"></a> -<span class="sourceLineNo">259</span> // If the mob file is empty, delete it instead of committing.<a name="line.259"></a> -<span class="sourceLineNo">260</span> store.getFileSystem().delete(mobFileWriter.getPath(), true);<a name="line.260"></a> -<span class="sourceLineNo">261</span> } catch (IOException e) {<a name="line.261"></a> -<span class="sourceLineNo">262</span> LOG.error("Failed to delete the temp mob file", e);<a name="line.262"></a> -<span class="sourceLineNo">263</span> }<a name="line.263"></a> -<span class="sourceLineNo">264</span> }<a name="line.264"></a> -<span class="sourceLineNo">265</span> }<a name="line.265"></a> -<span class="sourceLineNo">266</span>}<a name="line.266"></a> +<span class="sourceLineNo">237</span> }<a name="line.237"></a> +<span class="sourceLineNo">238</span><a name="line.238"></a> +<span class="sourceLineNo">239</span> if (mobCount > 0) {<a name="line.239"></a> +<span class="sourceLineNo">240</span> // commit the mob file from temp folder to target folder.<a name="line.240"></a> +<span class="sourceLineNo">241</span> // If the mob file is committed successfully but the store file is not,<a name="line.241"></a> +<span class="sourceLineNo">242</span> // the committed mob file will be handled by the sweep tool as an unused<a name="line.242"></a> +<span class="sourceLineNo">243</span> // file.<a name="line.243"></a> +<span class="sourceLineNo">244</span> status.setStatus("Flushing mob file " + store + ": appending metadata");<a name="line.244"></a> +<span class="sourceLineNo">245</span> mobFileWriter.appendMetadata(cacheFlushId, false, mobCount);<a name="line.245"></a> +<span class="sourceLineNo">246</span> status.setStatus("Flushing mob file " + store + ": closing flushed file");<a name="line.246"></a> +<span class="sourceLineNo">247</span> mobFileWriter.close();<a name="line.247"></a> +<span class="sourceLineNo">248</span> mobStore.commitFile(mobFileWriter.getPath(), targetPath);<a name="line.248"></a> +<span class="sourceLineNo">249</span> mobStore.updateMobFlushCount();<a name="line.249"></a> +<span class="sourceLineNo">250</span> mobStore.updateMobFlushedCellsCount(mobCount);<a name="line.250"></a> +<span class="sourceLineNo">251</span> mobStore.updateMobFlushedCellsSize(mobSize);<a name="line.251"></a> +<span class="sourceLineNo">252</span> } else {<a name="line.252"></a> +<span class="sourceLineNo">253</span> try {<a name="line.253"></a> +<span class="sourceLineNo">254</span> status.setStatus("Flushing mob file " + store + ": no mob cells, closing flushed file");<a name="line.254"></a> +<span class="sourceLineNo">255</span> mobFileWriter.close();<a name="line.255"></a> +<span class="sourceLineNo">256</span> // If the mob file is empty, delete it instead of committing.<a name="line.256"></a> +<span class="sourceLineNo">257</span> store.getFileSystem().delete(mobFileWriter.getPath(), true);<a name="line.257"></a> +<span class="sourceLineNo">258</span> } catch (IOException e) {<a name="line.258"></a> +<span class="sourceLineNo">259</span> LOG.error("Failed to delete the temp mob file", e);<a name="line.259"></a> +<span class="sourceLineNo">260</span> }<a name="line.260"></a> +<span class="sourceLineNo">261</span> }<a name="line.261"></a> +<span class="sourceLineNo">262</span> }<a name="line.262"></a> +<span class="sourceLineNo">263</span>}<a name="line.263"></a>
http://git-wip-us.apache.org/repos/asf/hbase-site/blob/9118853f/devapidocs/src-html/org/apache/hadoop/hbase/regionserver/DefaultStoreFlusher.html ---------------------------------------------------------------------- diff --git a/devapidocs/src-html/org/apache/hadoop/hbase/regionserver/DefaultStoreFlusher.html b/devapidocs/src-html/org/apache/hadoop/hbase/regionserver/DefaultStoreFlusher.html index b665347..e885f28 100644 --- a/devapidocs/src-html/org/apache/hadoop/hbase/regionserver/DefaultStoreFlusher.html +++ b/devapidocs/src-html/org/apache/hadoop/hbase/regionserver/DefaultStoreFlusher.html @@ -62,48 +62,44 @@ <span class="sourceLineNo">054</span> // Use a store scanner to find which rows to flush.<a name="line.54"></a> <span class="sourceLineNo">055</span> long smallestReadPoint = store.getSmallestReadPoint();<a name="line.55"></a> <span class="sourceLineNo">056</span> InternalScanner scanner = createScanner(snapshot.getScanners(), smallestReadPoint, tracker);<a name="line.56"></a> -<span class="sourceLineNo">057</span> if (scanner == null) {<a name="line.57"></a> -<span class="sourceLineNo">058</span> return result; // NULL scanner returned from coprocessor hooks means skip normal processing<a name="line.58"></a> -<span class="sourceLineNo">059</span> }<a name="line.59"></a> -<span class="sourceLineNo">060</span><a name="line.60"></a> -<span class="sourceLineNo">061</span> StoreFileWriter writer;<a name="line.61"></a> -<span class="sourceLineNo">062</span> try {<a name="line.62"></a> -<span class="sourceLineNo">063</span> // TODO: We can fail in the below block before we complete adding this flush to<a name="line.63"></a> -<span class="sourceLineNo">064</span> // list of store files. Add cleanup of anything put on filesystem if we fail.<a name="line.64"></a> -<span class="sourceLineNo">065</span> synchronized (flushLock) {<a name="line.65"></a> -<span class="sourceLineNo">066</span> status.setStatus("Flushing " + store + ": creating writer");<a name="line.66"></a> -<span class="sourceLineNo">067</span> // Write the map out to the disk<a name="line.67"></a> -<span class="sourceLineNo">068</span> writer = store.createWriterInTmp(cellsCount, store.getColumnFamilyDescriptor().getCompressionType(),<a name="line.68"></a> -<span class="sourceLineNo">069</span> /* isCompaction = */ false,<a name="line.69"></a> -<span class="sourceLineNo">070</span> /* includeMVCCReadpoint = */ true,<a name="line.70"></a> -<span class="sourceLineNo">071</span> /* includesTags = */ snapshot.isTagsPresent(),<a name="line.71"></a> -<span class="sourceLineNo">072</span> /* shouldDropBehind = */ false);<a name="line.72"></a> -<span class="sourceLineNo">073</span> IOException e = null;<a name="line.73"></a> -<span class="sourceLineNo">074</span> try {<a name="line.74"></a> -<span class="sourceLineNo">075</span> performFlush(scanner, writer, smallestReadPoint, throughputController);<a name="line.75"></a> -<span class="sourceLineNo">076</span> } catch (IOException ioe) {<a name="line.76"></a> -<span class="sourceLineNo">077</span> e = ioe;<a name="line.77"></a> -<span class="sourceLineNo">078</span> // throw the exception out<a name="line.78"></a> -<span class="sourceLineNo">079</span> throw ioe;<a name="line.79"></a> -<span class="sourceLineNo">080</span> } finally {<a name="line.80"></a> -<span class="sourceLineNo">081</span> if (e != null) {<a name="line.81"></a> -<span class="sourceLineNo">082</span> writer.close();<a name="line.82"></a> -<span class="sourceLineNo">083</span> } else {<a name="line.83"></a> -<span class="sourceLineNo">084</span> finalizeWriter(writer, cacheFlushId, status);<a name="line.84"></a> -<span class="sourceLineNo">085</span> }<a name="line.85"></a> -<span class="sourceLineNo">086</span> }<a name="line.86"></a> -<span class="sourceLineNo">087</span> }<a name="line.87"></a> -<span class="sourceLineNo">088</span> } finally {<a name="line.88"></a> -<span class="sourceLineNo">089</span> scanner.close();<a name="line.89"></a> -<span class="sourceLineNo">090</span> }<a name="line.90"></a> -<span class="sourceLineNo">091</span> LOG.info("Flushed, sequenceid=" + cacheFlushId +", memsize="<a name="line.91"></a> -<span class="sourceLineNo">092</span> + StringUtils.TraditionalBinaryPrefix.long2String(snapshot.getDataSize(), "", 1) +<a name="line.92"></a> -<span class="sourceLineNo">093</span> ", hasBloomFilter=" + writer.hasGeneralBloom() +<a name="line.93"></a> -<span class="sourceLineNo">094</span> ", into tmp file " + writer.getPath());<a name="line.94"></a> -<span class="sourceLineNo">095</span> result.add(writer.getPath());<a name="line.95"></a> -<span class="sourceLineNo">096</span> return result;<a name="line.96"></a> -<span class="sourceLineNo">097</span> }<a name="line.97"></a> -<span class="sourceLineNo">098</span>}<a name="line.98"></a> +<span class="sourceLineNo">057</span> StoreFileWriter writer;<a name="line.57"></a> +<span class="sourceLineNo">058</span> try {<a name="line.58"></a> +<span class="sourceLineNo">059</span> // TODO: We can fail in the below block before we complete adding this flush to<a name="line.59"></a> +<span class="sourceLineNo">060</span> // list of store files. Add cleanup of anything put on filesystem if we fail.<a name="line.60"></a> +<span class="sourceLineNo">061</span> synchronized (flushLock) {<a name="line.61"></a> +<span class="sourceLineNo">062</span> status.setStatus("Flushing " + store + ": creating writer");<a name="line.62"></a> +<span class="sourceLineNo">063</span> // Write the map out to the disk<a name="line.63"></a> +<span class="sourceLineNo">064</span> writer = store.createWriterInTmp(cellsCount, store.getColumnFamilyDescriptor().getCompressionType(),<a name="line.64"></a> +<span class="sourceLineNo">065</span> /* isCompaction = */ false,<a name="line.65"></a> +<span class="sourceLineNo">066</span> /* includeMVCCReadpoint = */ true,<a name="line.66"></a> +<span class="sourceLineNo">067</span> /* includesTags = */ snapshot.isTagsPresent(),<a name="line.67"></a> +<span class="sourceLineNo">068</span> /* shouldDropBehind = */ false);<a name="line.68"></a> +<span class="sourceLineNo">069</span> IOException e = null;<a name="line.69"></a> +<span class="sourceLineNo">070</span> try {<a name="line.70"></a> +<span class="sourceLineNo">071</span> performFlush(scanner, writer, smallestReadPoint, throughputController);<a name="line.71"></a> +<span class="sourceLineNo">072</span> } catch (IOException ioe) {<a name="line.72"></a> +<span class="sourceLineNo">073</span> e = ioe;<a name="line.73"></a> +<span class="sourceLineNo">074</span> // throw the exception out<a name="line.74"></a> +<span class="sourceLineNo">075</span> throw ioe;<a name="line.75"></a> +<span class="sourceLineNo">076</span> } finally {<a name="line.76"></a> +<span class="sourceLineNo">077</span> if (e != null) {<a name="line.77"></a> +<span class="sourceLineNo">078</span> writer.close();<a name="line.78"></a> +<span class="sourceLineNo">079</span> } else {<a name="line.79"></a> +<span class="sourceLineNo">080</span> finalizeWriter(writer, cacheFlushId, status);<a name="line.80"></a> +<span class="sourceLineNo">081</span> }<a name="line.81"></a> +<span class="sourceLineNo">082</span> }<a name="line.82"></a> +<span class="sourceLineNo">083</span> }<a name="line.83"></a> +<span class="sourceLineNo">084</span> } finally {<a name="line.84"></a> +<span class="sourceLineNo">085</span> scanner.close();<a name="line.85"></a> +<span class="sourceLineNo">086</span> }<a name="line.86"></a> +<span class="sourceLineNo">087</span> LOG.info("Flushed, sequenceid=" + cacheFlushId +", memsize="<a name="line.87"></a> +<span class="sourceLineNo">088</span> + StringUtils.TraditionalBinaryPrefix.long2String(snapshot.getDataSize(), "", 1) +<a name="line.88"></a> +<span class="sourceLineNo">089</span> ", hasBloomFilter=" + writer.hasGeneralBloom() +<a name="line.89"></a> +<span class="sourceLineNo">090</span> ", into tmp file " + writer.getPath());<a name="line.90"></a> +<span class="sourceLineNo">091</span> result.add(writer.getPath());<a name="line.91"></a> +<span class="sourceLineNo">092</span> return result;<a name="line.92"></a> +<span class="sourceLineNo">093</span> }<a name="line.93"></a> +<span class="sourceLineNo">094</span>}<a name="line.94"></a>
