gpt4 book ai didi

com.twitter.distributedlog.ZooKeeperClientBuilder类的使用及代码示例

转载 作者:知者 更新时间:2024-03-19 00:49:31 25 4
gpt4 key购买 nike

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

ZooKeeperClientBuilder介绍

[英]Builder to build zookeeper client.
[中]生成器来构建zookeeper客户端。

代码示例

代码示例来源:origin: twitter/distributedlog

URI uri = URI.create(args[0]);
ZooKeeperClient zkc = ZooKeeperClientBuilder.newBuilder().uri(uri)
    .sessionTimeoutMs(10000).build();
BKDLConfig bkdlConfig;
try {

代码示例来源:origin: twitter/distributedlog

private static ZooKeeperClientBuilder createDLZKClientBuilder(String zkcName,
                              DistributedLogConfiguration conf,
                              String zkServers,
                              StatsLogger statsLogger) {
  RetryPolicy retryPolicy = null;
  if (conf.getZKNumRetries() > 0) {
    retryPolicy = new BoundExponentialBackoffRetryPolicy(
      conf.getZKRetryBackoffStartMillis(),
      conf.getZKRetryBackoffMaxMillis(), conf.getZKNumRetries());
  }
  ZooKeeperClientBuilder builder = ZooKeeperClientBuilder.newBuilder()
    .name(zkcName)
    .sessionTimeoutMs(conf.getZKSessionTimeoutMilliseconds())
    .retryThreadCount(conf.getZKClientNumberRetryThreads())
    .requestRateLimit(conf.getZKRequestRateLimit())
    .zkServers(zkServers)
    .retryPolicy(retryPolicy)
    .statsLogger(statsLogger)
    .zkAclId(conf.getZkAclId());
  LOG.info("Created shared zooKeeper client builder {}: zkServers = {}, numRetries = {}, sessionTimeout = {}, retryBackoff = {},"
       + " maxRetryBackoff = {}, zkAclId = {}.", new Object[] { zkcName, zkServers, conf.getZKNumRetries(),
      conf.getZKSessionTimeoutMilliseconds(), conf.getZKRetryBackoffStartMillis(),
      conf.getZKRetryBackoffMaxMillis(), conf.getZkAclId() });
  return builder;
}

代码示例来源:origin: twitter/distributedlog

/**
 * Return a zookeeper client builder for testing.
 *
 * @return a zookeeper client builder
 */
public static ZooKeeperClientBuilder newBuilder() {
  return ZooKeeperClientBuilder.newBuilder()
      .retryPolicy(RetryPolicyUtils.DEFAULT_INFINITE_RETRY_POLICY)
      .connectionTimeoutMs(10000)
      .sessionTimeoutMs(60000)
      .zkAclId(null)
      .statsLogger(NullStatsLogger.INSTANCE);
}

代码示例来源:origin: twitter/distributedlog

/**
 * Run given <i>handler</i> by providing an available new zookeeper client
 *
 * @param handler
 *          Handler to process with provided zookeeper client.
 * @param conf
 *          Distributedlog Configuration.
 * @param namespace
 *          Distributedlog Namespace.
 */
private static <T> T withZooKeeperClient(ZooKeeperClientHandler<T> handler,
                     DistributedLogConfiguration conf,
                     URI namespace) throws IOException {
  ZooKeeperClient zkc = ZooKeeperClientBuilder.newBuilder()
      .name(String.format("dlzk:%s:factory_static", namespace))
      .sessionTimeoutMs(conf.getZKSessionTimeoutMilliseconds())
      .uri(namespace)
      .retryThreadCount(conf.getZKClientNumberRetryThreads())
      .requestRateLimit(conf.getZKRequestRateLimit())
      .zkAclId(conf.getZkAclId())
      .build();
  try {
    return handler.handle(zkc);
  } finally {
    zkc.close();
  }
}

代码示例来源:origin: twitter/distributedlog

private ZooKeeperClientBuilder clientBuilder(int sessionTimeoutMs)
    throws Exception {
  return ZooKeeperClientBuilder.newBuilder()
      .name("zkc")
      .uri(DLMTestUtil.createDLMURI(zkPort, "/"))
      .sessionTimeoutMs(sessionTimeoutMs)
      .zkServers(zkServers)
      .retryPolicy(new BoundExponentialBackoffRetryPolicy(100, 200, 2));
}

代码示例来源:origin: twitter/distributedlog

conf.getZKRetryBackoffMaxMillis(),
    Integer.MAX_VALUE);
ZooKeeperClient zkc = ZooKeeperClientBuilder.newBuilder()
    .name("DLAuditor-ZK")
    .zkServers(zkServers)
    .sessionTimeoutMs(conf.getZKSessionTimeoutMilliseconds())
    .retryPolicy(retryPolicy)
    .zkAclId(conf.getZkAclId())
    .build();
ExecutorService executorService = Executors.newCachedThreadPool();
try {

代码示例来源:origin: twitter/distributedlog

@Before
public void setup() throws Exception {
  zkc = TestZooKeeperClientBuilder.newBuilder()
      .uri(createURI("/"))
      .sessionTimeoutMs(10000)
      .build();
  resolver = new ZkMetadataResolver(zkc);
}

代码示例来源:origin: twitter/distributedlog

/**
   * Create a zookeeper client builder with provided <i>conf</i> for testing.
   *
   * @param conf distributedlog configuration
   * @return zookeeper client builder
   */
  public static ZooKeeperClientBuilder newBuilder(DistributedLogConfiguration conf) {
    return ZooKeeperClientBuilder.newBuilder()
        .retryPolicy(RetryPolicyUtils.DEFAULT_INFINITE_RETRY_POLICY)
        .sessionTimeoutMs(conf.getZKSessionTimeoutMilliseconds())
        .zkAclId(conf.getZkAclId())
        .retryThreadCount(conf.getZKClientNumberRetryThreads())
        .requestRateLimit(conf.getZKRequestRateLimit())
        .statsLogger(NullStatsLogger.INSTANCE);
  }
}

代码示例来源:origin: twitter/distributedlog

@Before
public void setup() throws Exception {
  zooKeeperClient =
    TestZooKeeperClientBuilder.newBuilder()
      .uri(createDLMURI("/"))
      .build();
}

代码示例来源:origin: twitter/distributedlog

@Test(timeout = 60000)
public void testAclPermsZkAccessNoConflict() throws Exception {
  String namespace = "/" + runtime.getMethodName();
  initDlogMeta(namespace, "test-un", "test-stream");
  URI uri = createDLMURI(namespace);
  ZooKeeperClient zkc = TestZooKeeperClientBuilder.newBuilder()
    .name("unpriv")
    .uri(uri)
    .build();
  zkc.get().getChildren(uri.getPath() + "/test-stream", false, new Stat());
  zkc.get().getData(uri.getPath() + "/test-stream", false, new Stat());
}

代码示例来源:origin: twitter/distributedlog

@Before
public void setup() throws Exception {
  zkc = TestZooKeeperClientBuilder.newBuilder()
      .name("zkc")
      .uri(DLMTestUtil.createDLMURI(zkPort, "/"))
      .sessionTimeoutMs(sessionTimeoutMs)
      .build();
}

代码示例来源:origin: twitter/distributedlog

private ZooKeeperClient buildClient() throws Exception {
  return clientBuilder().zkAclId(null).build();
}

代码示例来源:origin: twitter/distributedlog

@Before
public void setup() throws Exception {
  zkc = TestZooKeeperClientBuilder.newBuilder()
      .uri(createURI("/"))
      .zkServers(zkServers)
      .build();
  bkc = BookKeeperClientBuilder.newBuilder().name("bkc")
      .dlConfig(dlConf).ledgersPath(ledgersPath).zkc(zkc).build();
}

代码示例来源:origin: twitter/distributedlog

/**
 * {@link https://issues.apache.org/jira/browse/DL-34}
 */
@DistributedLogAnnotations.FlakyTest
@Ignore
@Test(timeout = 60000)
public void testAclAuthSpansExpirationNonRetryableClient() throws Exception {
  ZooKeeperClient zkcAuth = clientBuilder().retryPolicy(null).zkAclId("test").build();
  zkcAuth.get().create("/test", new byte[0], DistributedLogConstants.EVERYONE_READ_CREATOR_ALL, CreateMode.PERSISTENT);
  CountDownLatch expired = awaitConnectionEvent(KeeperState.Expired, zkcAuth);
  CountDownLatch connected = awaitConnectionEvent(KeeperState.SyncConnected, zkcAuth);
  expireZooKeeperSession(zkcAuth.get(), 2000);
  expired.await(2, TimeUnit.SECONDS);
  connected.await(2, TimeUnit.SECONDS);
  zkcAuth.get().create("/test/key1", new byte[0], DistributedLogConstants.EVERYONE_READ_CREATOR_ALL, CreateMode.PERSISTENT);
  rmAll(zkcAuth, "/test");
}

代码示例来源:origin: twitter/distributedlog

@Test(timeout = 60000, expected = BKException.ZKException.class)
public void testOpenLedgerWhenZkClosed() throws Exception {
  ZooKeeperClient newZkc = TestZooKeeperClientBuilder.newBuilder()
      .name("zkc-openledger-when-zk-closed")
      .zkServers(zkServers)
      .build();
  BookKeeperClient newBkc = BookKeeperClientBuilder.newBuilder()
      .name("bkc-openledger-when-zk-closed")
      .zkc(newZkc)
      .ledgersPath(ledgersPath)
      .dlConfig(conf)
      .build();
  try {
    LedgerHandle lh = newBkc.get().createLedger(BookKeeper.DigestType.CRC32, "zkcClosed".getBytes(UTF_8));
    lh.close();
    newZkc.close();
    LedgerHandleCache cache =
        LedgerHandleCache.newBuilder().bkc(newBkc).conf(conf).build();
    // open ledger after zkc closed
    cache.openLedger(new LogSegmentMetadata.LogSegmentMetadataBuilder("",
        2, lh.getId(), 1).setLogSegmentSequenceNo(lh.getId()).build(), false);
  } finally {
    newBkc.close();
  }
}

代码示例来源:origin: twitter/distributedlog

@Before
public void setup() throws Exception {
  zkc = TestZooKeeperClientBuilder.newBuilder()
      .zkServers(zkServers)
      .build();
}

代码示例来源:origin: twitter/distributedlog

DLUtils.getZKServersFromDLUri(_uri),
    _statsLogger.scope("dlzk_factory_writer_shared"));
ZooKeeperClient nsZkc = nsZkcBuilder.build();

代码示例来源:origin: twitter/distributedlog

public static void unbind(URI uri) throws IOException {
  DistributedLogConfiguration conf = new DistributedLogConfiguration();
  ZooKeeperClient zkc = ZooKeeperClientBuilder.newBuilder()
      .sessionTimeoutMs(conf.getZKSessionTimeoutMilliseconds())
      .retryThreadCount(conf.getZKClientNumberRetryThreads())
      .requestRateLimit(conf.getZKRequestRateLimit())
      .zkAclId(conf.getZkAclId())
      .uri(uri)
      .build();
  byte[] data = new byte[0];
  try {
    zkc.get().setData(uri.getPath(), data, -1);
  } catch (KeeperException ke) {
    throw new IOException("Fail to unbound dl metadata on uri " + uri, ke);
  } catch (InterruptedException ie) {
    throw new IOException("Interrupted when unbinding dl metadata on uri " + uri, ie);
  } finally {
    zkc.close();
  }
}

代码示例来源:origin: twitter/distributedlog

conf.getZKRetryBackoffMaxMillis(),
    Integer.MAX_VALUE);
ZooKeeperClient zkc = ZooKeeperClientBuilder.newBuilder()
    .name("DLAuditor-ZK")
    .zkServers(zkServers)
    .sessionTimeoutMs(conf.getZKSessionTimeoutMilliseconds())
    .retryPolicy(retryPolicy)
    .zkAclId(conf.getZkAclId())
    .build();
ExecutorService executorService = Executors.newCachedThreadPool();
try {

代码示例来源:origin: twitter/distributedlog

@Before
public void setup() throws Exception {
  zkc = TestZooKeeperClientBuilder.newBuilder()
      .uri(createDLMURI("/"))
      .sessionTimeoutMs(zkSessionTimeoutMs)
      .build();
  scheduler = OrderedScheduler.newBuilder()
      .name("test-zk-namespace-watcher")
      .corePoolSize(1)
      .build();
}

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