MemStore刷写线程—MemStoreFlusher源代码分析,hbasememstore
在HBase中表由一个或多个Region组成,而Region由一个或者多个Store组成,Store又由一个MenStore和若干个StoreFile组成。无论是向HBase写入数据还是请求读数据,都首先经过MemStore,对于写请求来说就是将数据直接写入MemStore,对于读请求来说就是先检查MenStore中是否包含相应的数据,如果有则直接读取该数据,否则在StoreFile中检索并读取数据。当写入MemStore中的数据达到一定的数量时,就需要将其中的数据刷写到StoreFile中,这是由后台线程自动完成的,负责该任务的主要Java类是MemStoreFlusher,本篇文章将结合该类的源代码学习HBase如何决定是否刷写数据到StoreFile中的,而不关注MemStoreFlusher如何被实例化及启动的(相应的代码位于HRegionServer.java中)以及写入StoreFile的代码,而仅仅关注HBase在满足什么条件的情况将触发刷写。
官方文档将MenStoreFlusher解释为一个线程,而实际上该类既没有继承自Thread也没有实现Runnable接口,但定义了FlushHandler内部类,该类继承自HasThread,而HasThread实现了Runnable接口,在稍后的学习中将会看到FlushHandler就是负责刷写数据到StoreFile中的线程,并会根据参数hbase.hstore.flusher.count的值(默认为2)决定启动该线程的数量。该参数具有比较重要的意义,如果该值较小,则会导致刷写队列中包含过多的待刷写的Region,而如果该值较大的话,则会有较多并行执行刷写的线程,则会增加HDFS的负担,进而引起HBase更加频繁地执行Compaction。在MenStoreFlusher中刷写队列是由下面的数据结构定义的:
private final BlockingQueue<FlushQueueEntry> flushQueue = new DelayQueue<FlushQueueEntry>(); private final Map<HRegion, FlushRegionEntry> regionsInQueue = new HashMap<HRegion, FlushRegionEntry>();
这两个数据结构必须一起使用,若某个FlushQueueEntry在二者中的一个,则在另一个中也必须能够找到该FlushQueueEntry。注意FlushRegionEntry实现了FlushQueueEntry接口,而后者继承自java.util.concurrent.Delayed接口。另一个实现了FlushQueueEntry的类为WakeupFlushThread类,该类的主要用作占位符插入到刷写队列中以确保刷写线程不会休眠。FlushRegionEntry类保存了请求刷写的Region、重试次数以及在刷写队列中存在的时间。刷写队列flushQueue的类型为DelayQueue,保存在该队列中的对象必须实现了Delayed接口,且位于该队列中的对象只有在其延迟过期后才可以被取出,位于队列头部的是延迟过期最长的对象,如果没有延迟过期,队列没有头部,poll操作将返回null。当调用对象的getDelay(TimeUnit.NANOSECONDS)方法返回值小于等于0时,该对象的延迟过期。由于DelayQueue实现了BlockingQueue接口,该接口是线程安全的,因此该队列也就支持BlockingQueue具有的阻塞入列和出列,即当队列为空时,将等待队列直到队列不为空再提取对象,当队列为满时,将等待队列有空闲位置时再插入对象。FlushRegionEntry的该方法实现如下:
@Override
public long getDelay(TimeUnit unit) {
return unit.convert(this.whenToExpire - EnvironmentEdgeManager.currentTime(), TimeUnit.MILLISECONDS);
}
FlushRegionEntry除了实现getDelay方法外,还定义了requeue()方法:
public FlushRegionEntry requeue(final long when) {
this.whenToExpire = EnvironmentEdgeManager.currentTime() + when;
this.requeueCount++;
return this;
}
稍后会介绍为什么会使用DelayQueue这样的数据结构定义刷写队列,现在继续看MemStoreFlusher的源代码。在MemStoreFlusher的构造函数中,读取hbase-site.xml中设置的与刷写相关的参数并赋值给相应的变量,并实例化了刷写线程FlushHandler:
public MemStoreFlusher(final Configuration conf, final HRegionServer server) {
super();
this.server = server;
//hbase.server.thread.wakefrequency,默认值10s
this.threadWakeFrequency =conf.getLong(HConstants.THREAD_WAKE_FREQUENCY, 10 * 1000);
long max = ManagementFactory.getMemoryMXBean().getHeapMemoryUsage().getMax();
//hbase.regionserver.global.memstore.size的值,默认为0.4,且该值必须小于等于0.8,大于0,即(0.0,0.8]
float globalMemStorePercent = HeapMemorySizeUtil.getGlobalMemStorePercent(conf, true);
this.globalMemStoreLimit = (long) (max * globalMemStorePercent);
//hbase.regionserver.global.memstore.size.lower.limit
this.globalMemStoreLimitLowMarkPercent = HeapMemorySizeUtil.getGlobalMemStoreLowerMark(conf,globalMemStorePercent);
//RS Xmx * hbase.regionserver.global.memstore.size * hbase.regionserver.global.memstore.size.lower.limit
this.globalMemStoreLimitLowMark = (long) (this.globalMemStoreLimit * this.globalMemStoreLimitLowMarkPercent);
//默认值为90s,如果任何一个Store总的StoreFile数量大于hbase.hstore.blockingStoreFiles值,HRegion将阻塞更新直到该参数
//设置的时间到达或者Compaction完成。在该参数设置的时间到期后,即使Compaction没有完成,HRegion也将不再阻塞更新
this.blockingWaitTime = conf.getInt("hbase.hstore.blockingWaitTime", 90000);
int handlerCount = conf.getInt("hbase.hstore.flusher.count", 2);
this.flushHandlers = new FlushHandler[handlerCount];
}
MemStoreFlusher定义了三个刷写方法,分别为:flushOneForGlobalPressure、flushRegion(final FlushRegionEntry fqe)和flushRegion(finalHRegion region, final boolean emergencyFlush)其中前两个方法最终调用了最后一个方法,也就是只有最后一个方法执行了实际的刷写操作,前两个方法都在FlushHandler刷写线程中调用,分别用于当RegionServer中的所有MemStore的大小超过了RS Xmx *hbase.regionserver.global.memstore.size *hbase.regionserver.global.memstore.size.lower.limit进行的刷写和当某个Region需要执行刷写操作。
首先看看flushOneForGlobalPressure方法的作用,从该方法的名称就可以看出,该方法用于缓解全局MemStore缓存的压力,在FlushHandler刷写线程的run方法中,当刷写队列中没有要求立刻执行刷写的

