HDFS Centrailzed Cache会放到哪个DataNode上

前几天,在看Hadoop User Email List的时候,发现了一个关于HDFS Centrailzed Cache的问题。刚好我又不熟悉这块,甚至之前都没听说过,就好好了解了一番。

其实原理很简单,各位读一下下面的几个链接,就能清楚是怎么回事:

其实上面第二篇文章中已经介绍了,会将Cache放在这个 block 的三个 replica 所在的 DataNode 其中的剩余可用内存最多的一个上。但是当时我没细看,就自己阅读了一下源码来探究这个问题。

相关代码主要在CacheReplicationMonitor.chooseDatanodesForCaching中:

  /**
   * Chooses datanode locations for caching from a list of valid possibilities.
   * Non-stale nodes are chosen before stale nodes.
   *
   * @param possibilities List of candidate datanodes
   * @param neededCached Number of replicas needed
   * @param staleInterval Age of a stale datanode
   * @return A list of chosen datanodes
   */
  private static List<DatanodeDescriptor> chooseDatanodesForCaching(
      final List<DatanodeDescriptor> possibilities, final int neededCached,
      final long staleInterval) {
    // Make a copy that we can modify
    List<DatanodeDescriptor> targets =
        new ArrayList<DatanodeDescriptor>(possibilities);
    // Selected targets
    List<DatanodeDescriptor> chosen = new LinkedList<DatanodeDescriptor>();

    // Filter out stale datanodes
    List<DatanodeDescriptor> stale = new LinkedList<DatanodeDescriptor>();
    Iterator<DatanodeDescriptor> it = targets.iterator();
    while (it.hasNext()) {
      DatanodeDescriptor d = it.next();
      if (d.isStale(staleInterval)) {
        it.remove();
        stale.add(d);
      }
    }
    // Select targets
    while (chosen.size() < neededCached) {
      // Try to use stale nodes if we're out of non-stale nodes, else we're done
      if (targets.isEmpty()) {
        if (!stale.isEmpty()) {
          targets = stale;
        } else {
          break;
        }
      }
      // Select a random target
      DatanodeDescriptor target =
          chooseRandomDatanodeByRemainingCapacity(targets);
      chosen.add(target);
      targets.remove(target);
    }
    return chosen;
  }

以及CacheReplicationMonitor.addNewPendingCached中:

 /**
  * Add new entries to the PendingCached list.
  *
  * @param neededCached     The number of replicas that need to be cached.
  * @param cachedBlock      The block which needs to be cached.
  * @param cached           A list of DataNodes currently caching the block.
  * @param pendingCached    A list of DataNodes that will soon cache the
  *                         block.
  */
 private void addNewPendingCached(final int neededCached,
     CachedBlock cachedBlock, List<DatanodeDescriptor> cached,
     List<DatanodeDescriptor> pendingCached) {
   // To figure out which replicas can be cached, we consult the
   // blocksMap.  We don't want to try to cache a corrupt replica, though.
   BlockInfoContiguous blockInfo = blockManager.
         getStoredBlock(new Block(cachedBlock.getBlockId()));
   if (blockInfo == null) {
     LOG.debug("Block {}: can't add new cached replicas," +
         " because there is no record of this block " +
         "on the NameNode.", cachedBlock.getBlockId());
     return;
   }
   if (!blockInfo.isComplete()) {
     LOG.debug("Block {}: can't cache this block, because it is not yet"
         + " complete.", cachedBlock.getBlockId());
     return;
   }
   // Filter the list of replicas to only the valid targets
   List<DatanodeDescriptor> possibilities =
       new LinkedList<DatanodeDescriptor>();
   int numReplicas = blockInfo.getCapacity();
   Collection<DatanodeDescriptor> corrupt =
       blockManager.getCorruptReplicas(blockInfo);
   int outOfCapacity = 0;
   for (int i = 0; i < numReplicas; i++) {
     DatanodeDescriptor datanode = blockInfo.getDatanode(i);
     if (datanode == null) {
       continue;
     }
     if (datanode.isDecommissioned() || datanode.isDecommissionInProgress()) {
       continue;
     }
     if (corrupt != null && corrupt.contains(datanode)) {
       continue;
     }
     if (pendingCached.contains(datanode) || cached.contains(datanode)) {
       continue;
     }
     long pendingBytes = 0;
     // Subtract pending cached blocks from effective capacity
     Iterator<CachedBlock> it = datanode.getPendingCached().iterator();
     while (it.hasNext()) {
       CachedBlock cBlock = it.next();
       BlockInfoContiguous info =
           blockManager.getStoredBlock(new Block(cBlock.getBlockId()));
       if (info != null) {
         pendingBytes -= info.getNumBytes();
       }
     }
     it = datanode.getPendingUncached().iterator();
     // Add pending uncached blocks from effective capacity
     while (it.hasNext()) {
       CachedBlock cBlock = it.next();
       BlockInfoContiguous info =
           blockManager.getStoredBlock(new Block(cBlock.getBlockId()));
       if (info != null) {
         pendingBytes += info.getNumBytes();
       }
     }
     long pendingCapacity = pendingBytes + datanode.getCacheRemaining();
     if (pendingCapacity < blockInfo.getNumBytes()) {
       LOG.trace("Block {}: DataNode {} is not a valid possibility " +
           "because the block has size {}, but the DataNode only has {}" +
           "bytes of cache remaining ({} pending bytes, {} already cached.",
           blockInfo.getBlockId(), datanode.getDatanodeUuid(),
           blockInfo.getNumBytes(), pendingCapacity, pendingBytes,
           datanode.getCacheRemaining());
       outOfCapacity++;
       continue;
     }
     possibilities.add(datanode);
   }
   List<DatanodeDescriptor> chosen = chooseDatanodesForCaching(possibilities,
       neededCached, blockManager.getDatanodeManager().getStaleInterval());
   for (DatanodeDescriptor datanode : chosen) {
     LOG.trace("Block {}: added to PENDING_CACHED on DataNode {}",
         blockInfo.getBlockId(), datanode.getDatanodeUuid());
     pendingCached.add(datanode);
     boolean added = datanode.getPendingCached().add(cachedBlock);
     assert added;
   }
   // We were unable to satisfy the requested replication factor
   if (neededCached > chosen.size()) {
     LOG.debug("Block {}: we only have {} of {} cached replicas."
             + " {} DataNodes have insufficient cache capacity.",
         blockInfo.getBlockId(),
         (cachedBlock.getReplication() - neededCached + chosen.size()),
         cachedBlock.getReplication(), outOfCapacity
     );
   }
 }

我们可以看到,逻辑就是,从满足下面这几个条件的DataNode中,选择一个可用Cache内存最多的节点:

  1. 此DataNode状态正常,没有decommission
  2. 这个Block在这个DataNode上,不是corrupt状态
  3. 这个Block没有在DataNode上面缓存过(pendingCached以及cached中都没有此DataNode)
  4. 这个Block已经关闭(比如当一个Block在进行replication的时候,如果第二个replicate 2 和replication 3没有完成,那么就不能选择这些DataNode)

那策略介绍完了。这里我就有一个新的疑问了,我们知道,MapReduce的Mapper任务,会尽量被分配到有相应Block的那个节点上。那这儿会考虑Centrailzed Cache么? 对Spark等其它框架呢?

关于Mapper任务分配时是否会考虑Centrailzed Cache,我会查看相关源码,并整理上来。

©著作权归作者所有,转载或内容合作请联系作者
  • 序言:七十年代末,一起剥皮案震惊了整个滨河市,随后出现的几起案子,更是在滨河造成了极大的恐慌,老刑警刘岩,带你破解...
    沈念sama阅读 216,125评论 6 498
  • 序言:滨河连续发生了三起死亡事件,死亡现场离奇诡异,居然都是意外死亡,警方通过查阅死者的电脑和手机,发现死者居然都...
    沈念sama阅读 92,293评论 3 392
  • 文/潘晓璐 我一进店门,熙熙楼的掌柜王于贵愁眉苦脸地迎上来,“玉大人,你说我怎么就摊上这事。” “怎么了?”我有些...
    开封第一讲书人阅读 162,054评论 0 351
  • 文/不坏的土叔 我叫张陵,是天一观的道长。 经常有香客问我,道长,这世上最难降的妖魔是什么? 我笑而不...
    开封第一讲书人阅读 58,077评论 1 291
  • 正文 为了忘掉前任,我火速办了婚礼,结果婚礼上,老公的妹妹穿的比我还像新娘。我一直安慰自己,他们只是感情好,可当我...
    茶点故事阅读 67,096评论 6 388
  • 文/花漫 我一把揭开白布。 她就那样静静地躺着,像睡着了一般。 火红的嫁衣衬着肌肤如雪。 梳的纹丝不乱的头发上,一...
    开封第一讲书人阅读 51,062评论 1 295
  • 那天,我揣着相机与录音,去河边找鬼。 笑死,一个胖子当着我的面吹牛,可吹牛的内容都是我干的。 我是一名探鬼主播,决...
    沈念sama阅读 39,988评论 3 417
  • 文/苍兰香墨 我猛地睁开眼,长吁一口气:“原来是场噩梦啊……” “哼!你这毒妇竟也来了?” 一声冷哼从身侧响起,我...
    开封第一讲书人阅读 38,817评论 0 273
  • 序言:老挝万荣一对情侣失踪,失踪者是张志新(化名)和其女友刘颖,没想到半个月后,有当地人在树林里发现了一具尸体,经...
    沈念sama阅读 45,266评论 1 310
  • 正文 独居荒郊野岭守林人离奇死亡,尸身上长有42处带血的脓包…… 初始之章·张勋 以下内容为张勋视角 年9月15日...
    茶点故事阅读 37,486评论 2 331
  • 正文 我和宋清朗相恋三年,在试婚纱的时候发现自己被绿了。 大学时的朋友给我发了我未婚夫和他白月光在一起吃饭的照片。...
    茶点故事阅读 39,646评论 1 347
  • 序言:一个原本活蹦乱跳的男人离奇死亡,死状恐怖,灵堂内的尸体忽然破棺而出,到底是诈尸还是另有隐情,我是刑警宁泽,带...
    沈念sama阅读 35,375评论 5 342
  • 正文 年R本政府宣布,位于F岛的核电站,受9级特大地震影响,放射性物质发生泄漏。R本人自食恶果不足惜,却给世界环境...
    茶点故事阅读 40,974评论 3 325
  • 文/蒙蒙 一、第九天 我趴在偏房一处隐蔽的房顶上张望。 院中可真热闹,春花似锦、人声如沸。这庄子的主人今日做“春日...
    开封第一讲书人阅读 31,621评论 0 21
  • 文/苍兰香墨 我抬头看了看天上的太阳。三九已至,却和暖如春,着一层夹袄步出监牢的瞬间,已是汗流浃背。 一阵脚步声响...
    开封第一讲书人阅读 32,796评论 1 268
  • 我被黑心中介骗来泰国打工, 没想到刚下飞机就差点儿被人妖公主榨干…… 1. 我叫王不留,地道东北人。 一个月前我还...
    沈念sama阅读 47,642评论 2 368
  • 正文 我出身青楼,却偏偏与公主长得像,于是被迫代替她去往敌国和亲。 传闻我的和亲对象是个残疾皇子,可洞房花烛夜当晚...
    茶点故事阅读 44,538评论 2 352

推荐阅读更多精彩内容