com.yahoo.pulsar.zookeeper.ZooKeeperCache类的使用及代码示例

x33g5p2x  于2022-02-05 转载在 其他  
字(11.4k)|赞(0)|评价(0)|浏览(104)

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

ZooKeeperCache介绍

暂无

代码示例

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker-common

public ZooKeeper getZooKeeper() {
  return this.cache.getZooKeeper();
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

private void initZK() throws PulsarServerException {
  String[] paths = new String[] { MANAGED_LEDGER_ROOT, OWNER_INFO_ROOT, LOCAL_POLICIES_ROOT };
  // initialize the zk client with values
  try {
    ZooKeeper zk = cache.getZooKeeper();
    for (String path : paths) {
      if (cache.exists(path)) {
        continue;
      }
      try {
        ZkUtils.createFullPathOptimistic(zk, path, new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
      } catch (KeeperException.NodeExistsException e) {
        // Ok
      }
    }
  } catch (Exception e) {
    LOG.error(e.getMessage(), e);
    throw new PulsarServerException(e);
  }
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

private boolean isBrokerActive(String candidateBroker) throws KeeperException, InterruptedException {
  Set<String> activeNativeBrokers = pulsar.getLocalZkCache().getChildren(LoadManager.LOADBALANCE_BROKERS_ROOT);
  for (String brokerHostPort : activeNativeBrokers) {
    if (candidateBroker.equals("http://" + brokerHostPort)) {
      if (LOG.isDebugEnabled()) {
        LOG.debug("Broker {} found for SLA Monitoring Namespace", brokerHostPort);
      }
      return true;
    }
  }
  if (LOG.isDebugEnabled()) {
    LOG.debug("Broker not found for SLA Monitoring Namespace {}",
        candidateBroker + ":" + config.getWebServicePort());
  }
  return false;
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

@POST
@Path("/{cluster}")
@ApiOperation(value = "Update the configuration for a cluster.", notes = "This operation requires Pulsar super-user privileges.")
@ApiResponses(value = { @ApiResponse(code = 204, message = "Cluster has been updated"),
    @ApiResponse(code = 403, message = "Don't have admin permission"),
    @ApiResponse(code = 404, message = "Cluster doesn't exist") })
public void updateCluster(@PathParam("cluster") String cluster, ClusterData clusterData) {
  validateSuperUserAccess();
  validatePoliciesReadOnlyAccess();
  try {
    String clusterPath = path("clusters", cluster);
    globalZk().setData(clusterPath, jsonMapper().writeValueAsBytes(clusterData), -1);
    globalZkCache().invalidate(clusterPath);
    log.info("[{}] Updated cluster {}", clientAppId(), cluster);
  } catch (KeeperException.NoNodeException e) {
    log.warn("[{}] Failed to update cluster {}: Does not exist", clientAppId(), cluster);
    throw new RestException(Status.NOT_FOUND, "Cluster does not exist");
  } catch (Exception e) {
    log.error("[{}] Failed to update cluster {}", clientAppId(), cluster, e);
    throw new RestException(e);
  }
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

/**
 * If load balancing is enabled, load shedding is enabled by default unless forced off by setting a flag in global
 * zk /admin/flags/load-shedding-unload-disabled
 *
 * @return false by default, unload is allowed in load shedding true if zk flag is set, unload is disabled
 */
public static boolean isUnloadDisabledInLoadShedding(final PulsarService pulsar) {
  if (!pulsar.getConfiguration().isLoadBalancerEnabled()) {
    return true;
  }
  boolean unloadDisabledInLoadShedding = false;
  try {
    unloadDisabledInLoadShedding = pulsar.getGlobalZkCache()
        .exists(AdminResource.LOAD_SHEDDING_UNLOAD_DISABLED_FLAG_PATH);
  } catch (Exception e) {
    log.warn("Unable to fetch contents of [{}] from global zookeeper",
        AdminResource.LOAD_SHEDDING_UNLOAD_DISABLED_FLAG_PATH, e);
  }
  return unloadDisabledInLoadShedding;
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

private static CompletableFuture<PartitionedTopicMetadata> fetchPartitionedTopicMetadataAsync(PulsarService pulsar,
    String path) {
  CompletableFuture<PartitionedTopicMetadata> metadataFuture = new CompletableFuture<>();
  try {
    // gets the number of partitions from the zk cache
    pulsar.getGlobalZkCache().getDataAsync(path, new Deserializer<PartitionedTopicMetadata>() {
      @Override
      public PartitionedTopicMetadata deserialize(String key, byte[] content) throws Exception {
        return jsonMapper().readValue(content, PartitionedTopicMetadata.class);
      }
    }).thenAccept(metadata -> {
      // if the partitioned topic is not found in zk, then the topic is not partitioned
      if (metadata.isPresent()) {
        metadataFuture.complete(metadata.get());
      } else {
        metadataFuture.complete(new PartitionedTopicMetadata());
      }
    }).exceptionally(ex -> {
      metadataFuture.completeExceptionally(ex);
      return null;
    });
  } catch (Exception e) {
    metadataFuture.completeExceptionally(e);
  }
  return metadataFuture;
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

public void close() {
    if (this.rackawarePolicyZkCache != null) {
      this.rackawarePolicyZkCache.stop();
    }
    if (this.clientIsolationZkCache != null) {
      this.clientIsolationZkCache.stop();
    }
  }
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

@Override
public void doFilter(ServletRequest req, ServletResponse resp, FilterChain chain)
    throws IOException, ServletException {
  try {
    String minApiVersion = pulsar.getLocalZkCache().getData(MIN_API_VERSION_PATH,
        Deserializers.STRING_DESERIALIZER).orElseThrow(() -> new KeeperException.NoNodeException());
    String requestApiVersion = getRequestApiVersion(req);
    if (shouldAllowRequest(req.getRemoteAddr(), minApiVersion, requestApiVersion)) {
      // Allow the request to continue by invoking the next filter in
      // the chain.
      chain.doFilter(req, resp);
    } else {
      // The client's API version is less than the min supported,
      // reject the request.
      HttpServletResponse httpResponse = (HttpServletResponse) resp;
      HttpServletResponseWrapper respWrapper = new HttpServletResponseWrapper(httpResponse);
      respWrapper.sendError(HttpServletResponse.SC_BAD_REQUEST, "Unsuported Client version");
    }
  } catch (Exception ex) {
    LOG.warn("[{}] Unable to safely determine client version eligibility. Allowing request",
        req.getRemoteAddr());
    chain.doFilter(req, resp);
  }
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

globalZkCache().invalidate(namespacePath);

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

private void setDynamicConfigurationToZK(String zkPath, Map<String, String> settings) throws IOException {
  byte[] settingBytes = ObjectMapperFactory.getThreadLocal().writeValueAsBytes(settings);
  try {
    if (pulsar.getLocalZkCache().exists(zkPath)) {
      pulsar.getZkClient().setData(zkPath, settingBytes, -1);
    } else {
      ZkUtils.createFullPathOptimistic(pulsar.getZkClient(), zkPath, settingBytes, Ids.OPEN_ACL_UNSAFE,
          CreateMode.PERSISTENT);
    }
  } catch (Exception e) {
    log.warn("Got exception when writing to ZooKeeper path [{}]:", zkPath, e);
  }
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

protected CompletableFuture<LookupResult> createLookupResult(String candidateBroker) throws Exception {
  CompletableFuture<LookupResult> lookupFuture = new CompletableFuture<>();
  try {
    checkArgument(StringUtils.isNotBlank(candidateBroker), "Lookup broker can't be null " + candidateBroker);
    URI uri = new URI(candidateBroker);
    String path = String.format("%s/%s:%s", LoadManager.LOADBALANCE_BROKERS_ROOT, uri.getHost(),
        uri.getPort());
    pulsar.getLocalZkCache().getDataAsync(path, pulsar.getLoadManager().get().getLoadReportDeserializer()).thenAccept(reportData -> {
      if (reportData.isPresent()) {
        ServiceLookupData lookupData = reportData.get();
        lookupFuture.complete(new LookupResult(lookupData.getWebServiceUrl(),
            lookupData.getWebServiceUrlTls(), lookupData.getPulsarServiceUrl(),
            lookupData.getPulsarServiceUrlTls()));
      } else {
        lookupFuture.completeExceptionally(new KeeperException.NoNodeException(path));
      }
    }).exceptionally(ex -> {
      lookupFuture.completeExceptionally(ex);
      return null;
    });
  } catch (Exception e) {
    lookupFuture.completeExceptionally(e);
  }
  return lookupFuture;
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

private void saveQuotaToZnode(String zpath, ResourceQuota quota) throws Exception {
  ZooKeeper zk = this.localZkCache.getZooKeeper();
  if (zk.exists(zpath, false) == null) {
    try {
      ZkUtils.createFullPathOptimistic(zk, zpath, new byte[0], Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
    } catch (KeeperException.NodeExistsException e) {
    }
  }
  zk.setData(zpath, this.jsonMapper.writeValueAsBytes(quota), -1);
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

/**
 * Checks whether the broker is allowed to do read-write operations based on the existence of a node in global
 * zookeeper.
 *
 * @throws WebApplicationException
 *             if broker has a read only access if broker is not connected to the global zookeeper
 */
public void validatePoliciesReadOnlyAccess() {
  boolean arePoliciesReadOnly = true;
  try {
    arePoliciesReadOnly = globalZkCache().exists(POLICIES_READONLY_FLAG_PATH);
  } catch (Exception e) {
    log.warn("Unable to fetch contents of [{}] from global zookeeper", POLICIES_READONLY_FLAG_PATH, e);
    throw new RestException(e);
  }
  if (arePoliciesReadOnly) {
    log.debug("Policies are read-only. Broker cannot do read-write operations");
    throw new RestException(Status.FORBIDDEN, "Broker is forbidden to do read-write operations");
  } else {
    // Make sure the broker is connected to the global zookeeper before writing. If not, throw an exception.
    if (globalZkCache().getZooKeeper().getState() != States.CONNECTED) {
      log.debug("Broker is not connected to the global zookeeper");
      throw new RestException(Status.PRECONDITION_FAILED,
          "Broker needs to be connected to global zookeeper before making a read-write operation");
    } else {
      // Do nothing, just log the message.
      log.debug("Broker is allowed to make read-write operations");
    }
  }
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

globalZkCache().invalidate(propertyPath);
  log.info("[{}] updated property {}", clientAppId(), property);
} catch (RestException re) {

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

@GET
@Path("/{cluster}")
@ApiOperation(value = "Get the list of active brokers (web service addresses) in the cluster.", response = String.class, responseContainer = "Set")
@ApiResponses(value = { @ApiResponse(code = 403, message = "Don't have admin permission"),
    @ApiResponse(code = 404, message = "Cluster doesn't exist") })
public Set<String> getActiveBrokers(@PathParam("cluster") String cluster) throws Exception {
  validateSuperUserAccess();
  validateClusterOwnership(cluster);
  try {
    // Add Native brokers
    return pulsar().getLocalZkCache().getChildren(LoadManager.LOADBALANCE_BROKERS_ROOT);
  } catch (Exception e) {
    LOG.error(String.format("[%s] Failed to get active broker list: cluster=%s", clientAppId(), cluster), e);
    throw new RestException(e);
  }
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker-common

private void initZK() throws PulsarServerException {
  String[] paths = new String[] { CLUSTERS_ROOT, POLICIES_ROOT };
  // initialize the zk client with values
  try {
    ZooKeeper zk = cache.getZooKeeper();
    for (String path : paths) {
      try {
        if (zk.exists(path, false) == null) {
          ZkUtils.createFullPathOptimistic(zk, path, new byte[0], Ids.OPEN_ACL_UNSAFE,
              CreateMode.PERSISTENT);
        }
      } catch (KeeperException.NodeExistsException e) {
        // Ok
      }
    }
  } catch (Exception e) {
    LOG.error(e.getMessage(), e);
    throw new PulsarServerException(e);
  }
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

try {
  globalZk().delete(path, -1);
  globalZkCache().invalidate(path);

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

protected ZooKeeper globalZk() {
  return pulsar().getGlobalZkCache().getZooKeeper();
}

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

String clusterPath = path("clusters", cluster);
  globalZk().delete(clusterPath, -1);
  globalZkCache().invalidate(clusterPath);
  log.info("[{}] Deleted cluster {}", clientAppId(), cluster);
} catch (KeeperException.NoNodeException e) {

代码示例来源:origin: com.yahoo.pulsar/pulsar-broker

/**
 * Disable bundle in local cache and on zk
 * 
 * @param bundle
 * @throws Exception
 */
public void disableOwnership(NamespaceBundle bundle) throws Exception {
  String path = ServiceUnitZkUtils.path(bundle);
  updateBundleState(bundle, false);
  localZkCache.getZooKeeper().setData(path, jsonMapper.writeValueAsBytes(selfOwnerInfoDisabled), -1);
  ownershipReadOnlyCache.invalidate(path);
}

相关文章

ZooKeeperCache类方法