gpt4 book ai didi

org.apache.hadoop.hbase.procedure.ZKProcedureUtil.getWatcher()方法的使用及代码示例

转载 作者:知者 更新时间:2024-03-17 02:25:31 29 4
gpt4 key购买 nike

本文整理了Java中org.apache.hadoop.hbase.procedure.ZKProcedureUtil.getWatcher()方法的一些代码示例,展示了ZKProcedureUtil.getWatcher()的具体用法。这些代码示例主要来源于Github/Stackoverflow/Maven等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。ZKProcedureUtil.getWatcher()方法的具体详情如下:
包路径:org.apache.hadoop.hbase.procedure.ZKProcedureUtil
类名称:ZKProcedureUtil
方法名:getWatcher

ZKProcedureUtil.getWatcher介绍

暂无

代码示例

代码示例来源:origin: apache/hbase

try {
 ZKUtil.createWithParents(zkProc.getWatcher(), reachedNode);
  if (ZKUtil.watchAndCheckExists(zkProc.getWatcher(), znode)) {
   byte[] dataFromMember = ZKUtil.getData(zkProc.getWatcher(), znode);

代码示例来源:origin: apache/hbase

try {
 if (ZKUtil.watchAndCheckExists(zkProc.getWatcher(), abortNode)) {
  abort(abortNode);
 ZKUtil.createWithParents(zkProc.getWatcher(), acquire, data);
  if (ZKUtil.watchAndCheckExists(zkProc.getWatcher(), znode)) {
   coordinator.memberAcquiredBarrier(procName, node);

代码示例来源:origin: apache/hbase

/**
 * This attempts to create an acquired state znode for the procedure (snapshot name).
 *
 * It then looks for the reached znode to trigger in-barrier execution.  If not present we
 * have a watcher, if present then trigger the in-barrier action.
 */
@Override
public void sendMemberAcquired(Subprocedure sub) throws IOException {
 String procName = sub.getName();
 try {
  LOG.debug("Member: '" + memberName + "' joining acquired barrier for procedure (" + procName
    + ") in zk");
  String acquiredZNode = ZNodePaths.joinZNode(ZKProcedureUtil.getAcquireBarrierNode(
   zkController, procName), memberName);
  ZKUtil.createAndFailSilent(zkController.getWatcher(), acquiredZNode);
  // watch for the complete node for this snapshot
  String reachedBarrier = zkController.getReachedBarrierNode(procName);
  LOG.debug("Watch for global barrier reached:" + reachedBarrier);
  if (ZKUtil.watchAndCheckExists(zkController.getWatcher(), reachedBarrier)) {
   receivedReachedGlobalBarrier(reachedBarrier);
  }
 } catch (KeeperException e) {
  member.controllerConnectionFailure("Failed to acquire barrier for procedure: "
    + procName + " and member: " + memberName, e, procName);
 }
}

代码示例来源:origin: apache/hbase

private void waitForNewProcedures() {
 // watch for new procedues that we need to start subprocedures for
 LOG.debug("Looking for new procedures under znode:'" + zkController.getAcquiredBarrier() + "'");
 List<String> runningProcedures = null;
 try {
  runningProcedures = ZKUtil.listChildrenAndWatchForNewChildren(zkController.getWatcher(),
   zkController.getAcquiredBarrier());
  if (runningProcedures == null) {
   LOG.debug("No running procedures.");
   return;
  }
 } catch (KeeperException e) {
  member.controllerConnectionFailure("General failure when watching for new procedures",
   e, null);
 }
 if (runningProcedures == null) {
  LOG.debug("No running procedures.");
  return;
 }
 for (String procName : runningProcedures) {
  // then read in the procedure information
  String path = ZNodePaths.joinZNode(zkController.getAcquiredBarrier(), procName);
  startNewSubprocedure(path);
 }
}

代码示例来源:origin: apache/hbase

/**
 * This is the abort message being sent by the coordinator to member
 *
 * TODO this code isn't actually used but can be used to issue a cancellation from the
 * coordinator.
 */
@Override
final public void sendAbortToMembers(Procedure proc, ForeignException ee) {
 String procName = proc.getName();
 LOG.debug("Aborting procedure '" + procName + "' in zk");
 String procAbortNode = zkProc.getAbortZNode(procName);
 try {
  LOG.debug("Creating abort znode:" + procAbortNode);
  String source = (ee.getSource() == null) ? coordName : ee.getSource();
  byte[] errorInfo = ProtobufUtil.prependPBMagic(ForeignException.serialize(source, ee));
  // first create the znode for the procedure
  ZKUtil.createAndFailSilent(zkProc.getWatcher(), procAbortNode, errorInfo);
  LOG.debug("Finished creating abort node:" + procAbortNode);
 } catch (KeeperException e) {
  // possible that we get this error for the procedure if we already reset the zk state, but in
  // that case we should still get an error for that procedure anyways
  zkProc.logZKTree(zkProc.baseZNode);
  coordinator.rpcConnectionFailure("Failed to post zk node:" + procAbortNode
    + " to abort procedure '" + procName + "'", new IOException(e));
 }
}

代码示例来源:origin: apache/hbase

private void watchForAbortedProcedures() {
 LOG.debug("Checking for aborted procedures on node: '" + zkController.getAbortZnode() + "'");
 try {
  // this is the list of the currently aborted procedues
  List<String> children = ZKUtil.listChildrenAndWatchForNewChildren(zkController.getWatcher(),
         zkController.getAbortZnode());
  if (children == null || children.isEmpty()) {
   return;
  }
  for (String node : children) {
   String abortNode = ZNodePaths.joinZNode(zkController.getAbortZnode(), node);
   abort(abortNode);
  }
 } catch (KeeperException e) {
  member.controllerConnectionFailure("Failed to list children for abort node:"
    + zkController.getAbortZnode(), e, null);
 }
}

代码示例来源:origin: apache/hbase

String opName = ZKUtil.getNodeName(abortZNode);
try {
 byte[] data = ZKUtil.getData(zkController.getWatcher(), abortZNode);

代码示例来源:origin: apache/hbase

if (ZKUtil.watchAndCheckExists(zkController.getWatcher(), abortZNode)) {
 LOG.debug("Not starting:" + opName + " because we already have an abort notification.");
 return;
byte[] data = ZKUtil.getData(zkController.getWatcher(), path);
if (!ProtobufUtil.isPBMagicPrefix(data)) {
 String msg = "Data in for starting procedure " + opName +

代码示例来源:origin: apache/hbase

ForeignException ee = null;
try {
 byte[] data = ZKUtil.getData(zkProc.getWatcher(), abortNode);
 if (data == null || data.length == 0) {

代码示例来源:origin: apache/hbase

/**
 * This acts as the ack for a completed procedure
 */
@Override
public void sendMemberCompleted(Subprocedure sub, byte[] data) throws IOException {
 String procName = sub.getName();
 LOG.debug("Marking procedure  '" + procName + "' completed for member '" + memberName
   + "' in zk");
 String joinPath =
  ZNodePaths.joinZNode(zkController.getReachedBarrierNode(procName), memberName);
 // ProtobufUtil.prependPBMagic does not take care of null
 if (data == null) {
  data = new byte[0];
 }
 try {
  ZKUtil.createAndFailSilent(zkController.getWatcher(), joinPath,
   ProtobufUtil.prependPBMagic(data));
 } catch (KeeperException e) {
  member.controllerConnectionFailure("Failed to post zk node:" + joinPath
    + " to join procedure barrier.", e, procName);
 }
}

代码示例来源:origin: apache/hbase

/**
 * This should be called by the member and should write a serialized root cause exception as
 * to the abort znode.
 */
@Override
public void sendMemberAborted(Subprocedure sub, ForeignException ee) {
 if (sub == null) {
  LOG.error("Failed due to null subprocedure", ee);
  return;
 }
 String procName = sub.getName();
 LOG.debug("Aborting procedure (" + procName + ") in zk");
 String procAbortZNode = zkController.getAbortZNode(procName);
 try {
  String source = (ee.getSource() == null) ? memberName: ee.getSource();
  byte[] errorInfo = ProtobufUtil.prependPBMagic(ForeignException.serialize(source, ee));
  ZKUtil.createAndFailSilent(zkController.getWatcher(), procAbortZNode, errorInfo);
  LOG.debug("Finished creating abort znode:" + procAbortZNode);
 } catch (KeeperException e) {
  // possible that we get this error for the procedure if we already reset the zk state, but in
  // that case we should still get an error for that procedure anyways
  zkController.logZKTree(zkController.getBaseZnode());
  member.controllerConnectionFailure("Failed to post zk node:" + procAbortZNode
    + " to abort procedure", e, procName);
 }
}

代码示例来源:origin: co.cask.hbase/hbase

@Override
public void sendGlobalBarrierReached(Procedure proc, List<String> nodeNames) throws IOException {
 String procName = proc.getName();
 String reachedNode = zkProc.getReachedBarrierNode(procName);
 LOG.debug("Creating reached barrier zk node:" + reachedNode);
 try {
  // create the reached znode and watch for the reached znodes
  ZKUtil.createWithParents(zkProc.getWatcher(), reachedNode);
  // loop through all the children of the acquire phase and watch for them
  for (String node : nodeNames) {
   String znode = ZKUtil.joinZNode(reachedNode, node);
   if (ZKUtil.watchAndCheckExists(zkProc.getWatcher(), znode)) {
    coordinator.memberFinishedBarrier(procName, node);
   }
  }
 } catch (KeeperException e) {
  throw new IOException("Failed while creating reached node:" + reachedNode, e);
 }
}

代码示例来源:origin: co.cask.hbase/hbase

/**
 * This acts as the ack for a completed snapshot
 */
@Override
public void sendMemberCompleted(Subprocedure sub) throws IOException {
 String procName = sub.getName();
 LOG.debug("Marking procedure  '" + procName + "' completed for member '" + memberName
   + "' in zk");
 String joinPath = ZKUtil.joinZNode(zkController.getReachedBarrierNode(procName), memberName);
 try {
  ZKUtil.createAndFailSilent(zkController.getWatcher(), joinPath);
 } catch (KeeperException e) {
  member.controllerConnectionFailure("Failed to post zk node:" + joinPath
    + " to join procedure barrier.", new IOException(e));
 }
}

代码示例来源:origin: co.cask.hbase/hbase

/**
 * This attempts to create an acquired state znode for the procedure (snapshot name).
 *
 * It then looks for the reached znode to trigger in-barrier execution.  If not present we
 * have a watcher, if present then trigger the in-barrier action.
 */
@Override
public void sendMemberAcquired(Subprocedure sub) throws IOException {
 String procName = sub.getName();
 try {
  LOG.debug("Member: '" + memberName + "' joining acquired barrier for procedure (" + procName
    + ") in zk");
  String acquiredZNode = ZKUtil.joinZNode(ZKProcedureUtil.getAcquireBarrierNode(
   zkController, procName), memberName);
  ZKUtil.createAndFailSilent(zkController.getWatcher(), acquiredZNode);
  // watch for the complete node for this snapshot
  String reachedBarrier = zkController.getReachedBarrierNode(procName);
  LOG.debug("Watch for global barrier reached:" + reachedBarrier);
  if (ZKUtil.watchAndCheckExists(zkController.getWatcher(), reachedBarrier)) {
   receivedReachedGlobalBarrier(reachedBarrier);
  }
 } catch (KeeperException e) {
  member.controllerConnectionFailure("Failed to acquire barrier for procedure: "
    + procName + " and member: " + memberName, new IOException(e));
 }
}

代码示例来源:origin: co.cask.hbase/hbase

private void waitForNewProcedures() {
 // watch for new procedues that we need to start subprocedures for
 LOG.debug("Looking for new procedures under znode:'" + zkController.getAcquiredBarrier() + "'");
 List<String> runningProcedures = null;
 try {
  runningProcedures = ZKUtil.listChildrenAndWatchForNewChildren(zkController.getWatcher(),
   zkController.getAcquiredBarrier());
  if (runningProcedures == null) {
   LOG.debug("No running procedures.");
   return;
  }
 } catch (KeeperException e) {
  member.controllerConnectionFailure("General failure when watching for new procedures",
   new IOException(e));
 }
 if (runningProcedures == null) {
  LOG.debug("No running procedures.");
  return;
 }
 for (String procName : runningProcedures) {
  // then read in the procedure information
  String path = ZKUtil.joinZNode(zkController.getAcquiredBarrier(), procName);
  startNewSubprocedure(path);
 }
}

代码示例来源:origin: co.cask.hbase/hbase

private void watchForAbortedProcedures() {
 LOG.debug("Checking for aborted procedures on node: '" + zkController.getAbortZnode() + "'");
 try {
  // this is the list of the currently aborted procedues
  for (String node : ZKUtil.listChildrenAndWatchForNewChildren(zkController.getWatcher(),
   zkController.getAbortZnode())) {
   String abortNode = ZKUtil.joinZNode(zkController.getAbortZnode(), node);
   abort(abortNode);
  }
 } catch (KeeperException e) {
  member.controllerConnectionFailure("Failed to list children for abort node:"
    + zkController.getAbortZnode(), new IOException(e));
 }
}

代码示例来源:origin: harbby/presto-connectors

/**
 * This acts as the ack for a completed procedure
 */
@Override
public void sendMemberCompleted(Subprocedure sub, byte[] data) throws IOException {
 String procName = sub.getName();
 LOG.debug("Marking procedure  '" + procName + "' completed for member '" + memberName
   + "' in zk");
 String joinPath = ZKUtil.joinZNode(zkController.getReachedBarrierNode(procName), memberName);
 // ProtobufUtil.prependPBMagic does not take care of null
 if (data == null) {
  data = new byte[0];
 }
 try {
  ZKUtil.createAndFailSilent(zkController.getWatcher(), joinPath,
   ProtobufUtil.prependPBMagic(data));
 } catch (KeeperException e) {
  member.controllerConnectionFailure("Failed to post zk node:" + joinPath
    + " to join procedure barrier.", e, procName);
 }
}

代码示例来源:origin: harbby/presto-connectors

private void watchForAbortedProcedures() {
 LOG.debug("Checking for aborted procedures on node: '" + zkController.getAbortZnode() + "'");
 try {
  // this is the list of the currently aborted procedues
  for (String node : ZKUtil.listChildrenAndWatchForNewChildren(zkController.getWatcher(),
   zkController.getAbortZnode())) {
   String abortNode = ZKUtil.joinZNode(zkController.getAbortZnode(), node);
   abort(abortNode);
  }
 } catch (KeeperException e) {
  member.controllerConnectionFailure("Failed to list children for abort node:"
    + zkController.getAbortZnode(), e, null);
 }
}

代码示例来源:origin: harbby/presto-connectors

/**
 * This is the abort message being sent by the coordinator to member
 *
 * TODO this code isn't actually used but can be used to issue a cancellation from the
 * coordinator.
 */
@Override
final public void sendAbortToMembers(Procedure proc, ForeignException ee) {
 String procName = proc.getName();
 LOG.debug("Aborting procedure '" + procName + "' in zk");
 String procAbortNode = zkProc.getAbortZNode(procName);
 try {
  LOG.debug("Creating abort znode:" + procAbortNode);
  String source = (ee.getSource() == null) ? coordName : ee.getSource();
  byte[] errorInfo = ProtobufUtil.prependPBMagic(ForeignException.serialize(source, ee));
  // first create the znode for the procedure
  ZKUtil.createAndFailSilent(zkProc.getWatcher(), procAbortNode, errorInfo);
  LOG.debug("Finished creating abort node:" + procAbortNode);
 } catch (KeeperException e) {
  // possible that we get this error for the procedure if we already reset the zk state, but in
  // that case we should still get an error for that procedure anyways
  zkProc.logZKTree(zkProc.baseZNode);
  coordinator.rpcConnectionFailure("Failed to post zk node:" + procAbortNode
    + " to abort procedure '" + procName + "'", new IOException(e));
 }
}

代码示例来源:origin: harbby/presto-connectors

private void waitForNewProcedures() {
 // watch for new procedues that we need to start subprocedures for
 LOG.debug("Looking for new procedures under znode:'" + zkController.getAcquiredBarrier() + "'");
 List<String> runningProcedures = null;
 try {
  runningProcedures = ZKUtil.listChildrenAndWatchForNewChildren(zkController.getWatcher(),
   zkController.getAcquiredBarrier());
  if (runningProcedures == null) {
   LOG.debug("No running procedures.");
   return;
  }
 } catch (KeeperException e) {
  member.controllerConnectionFailure("General failure when watching for new procedures",
   e, null);
 }
 if (runningProcedures == null) {
  LOG.debug("No running procedures.");
  return;
 }
 for (String procName : runningProcedures) {
  // then read in the procedure information
  String path = ZKUtil.joinZNode(zkController.getAcquiredBarrier(), procName);
  startNewSubprocedure(path);
 }
}

29 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com