org.apache.flink.runtime.highavailability.zookeeper.ZooKeeperRunningJobsRegistry类的使用及代码示例

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

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

ZooKeeperRunningJobsRegistry介绍

[英]A zookeeper based registry for running jobs, highly available.
[中]基于zookeeper的注册表,用于运行作业,高度可用。

代码示例

代码示例来源:origin: org.apache.flink/flink-runtime

@Override
public void setJobRunning(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    writeEnumToZooKeeper(jobID, JobSchedulingStatus.RUNNING);
  }
  catch (Exception e) {
    throw new IOException("Failed to set RUNNING state in ZooKeeper for job " + jobID, e);
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime_2.10

private void writeEnumToZooKeeper(JobID jobID, JobSchedulingStatus status) throws Exception {
    final String zkPath = createZkPath(jobID);
    this.client.newNamespaceAwareEnsurePath(zkPath).ensure(client.getZookeeperClient());
    this.client.setData().forPath(zkPath, status.name().getBytes(ENCODING));
  }
}

代码示例来源:origin: com.alibaba.blink/flink-runtime

public ZooKeeperHaServices(
    CuratorFramework client,
    Executor executor,
    Configuration configuration,
    BlobStoreService blobStoreService) {
  this.client = checkNotNull(client);
  this.executor = checkNotNull(executor);
  this.configuration = checkNotNull(configuration);
  this.runningJobsRegistry = new ZooKeeperRunningJobsRegistry(client, configuration);
  this.blobStoreService = checkNotNull(blobStoreService);
}

代码示例来源:origin: com.alibaba.blink/flink-runtime

@Override
public void clearJob(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    final String zkPath = createZkPath(jobID);
    this.client.newNamespaceAwareEnsurePath(zkPath).ensure(client.getZookeeperClient());
    this.client.delete().forPath(zkPath);
  }
  catch (Exception e) {
    throw new IOException("Failed to clear job state from ZooKeeper for job " + jobID, e);
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime_2.10

public ZooKeeperHaServices(
    CuratorFramework client,
    Executor executor,
    Configuration configuration,
    BlobStoreService blobStoreService) {
  this.client = checkNotNull(client);
  this.executor = checkNotNull(executor);
  this.configuration = checkNotNull(configuration);
  this.runningJobsRegistry = new ZooKeeperRunningJobsRegistry(client, configuration);
  this.blobStoreService = checkNotNull(blobStoreService);
}

代码示例来源:origin: org.apache.flink/flink-runtime_2.10

@Override
public void clearJob(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    final String zkPath = createZkPath(jobID);
    this.client.newNamespaceAwareEnsurePath(zkPath).ensure(client.getZookeeperClient());
    this.client.delete().forPath(zkPath);
  }
  catch (Exception e) {
    throw new IOException("Failed to clear job state from ZooKeeper for job " + jobID, e);
  }
}

代码示例来源:origin: com.alibaba.blink/flink-runtime

@Override
public void setJobRunning(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    writeEnumToZooKeeper(jobID, JobSchedulingStatus.RUNNING);
  }
  catch (Exception e) {
    throw new IOException("Failed to set RUNNING state in ZooKeeper for job " + jobID, e);
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime_2.11

public ZooKeeperHaServices(
    CuratorFramework client,
    Executor executor,
    Configuration configuration,
    BlobStoreService blobStoreService) {
  this.client = checkNotNull(client);
  this.executor = checkNotNull(executor);
  this.configuration = checkNotNull(configuration);
  this.runningJobsRegistry = new ZooKeeperRunningJobsRegistry(client, configuration);
  this.blobStoreService = checkNotNull(blobStoreService);
}

代码示例来源:origin: org.apache.flink/flink-runtime

private void writeEnumToZooKeeper(JobID jobID, JobSchedulingStatus status) throws Exception {
    final String zkPath = createZkPath(jobID);
    this.client.newNamespaceAwareEnsurePath(zkPath).ensure(client.getZookeeperClient());
    this.client.setData().forPath(zkPath, status.name().getBytes(ENCODING));
  }
}

代码示例来源:origin: com.alibaba.blink/flink-runtime

@Override
public void setJobFinished(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    writeEnumToZooKeeper(jobID, JobSchedulingStatus.DONE);
  }
  catch (Exception e) {
    throw new IOException("Failed to set DONE state in ZooKeeper for job " + jobID, e);
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime

public ZooKeeperHaServices(
    CuratorFramework client,
    Executor executor,
    Configuration configuration,
    BlobStoreService blobStoreService) {
  this.client = checkNotNull(client);
  this.executor = checkNotNull(executor);
  this.configuration = checkNotNull(configuration);
  this.runningJobsRegistry = new ZooKeeperRunningJobsRegistry(client, configuration);
  this.blobStoreService = checkNotNull(blobStoreService);
}

代码示例来源:origin: org.apache.flink/flink-runtime_2.11

@Override
public JobSchedulingStatus getJobSchedulingStatus(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    final String zkPath = createZkPath(jobID);
    final Stat stat = client.checkExists().forPath(zkPath);
    if (stat != null) {
      // found some data, try to parse it
      final byte[] data = client.getData().forPath(zkPath);
      if (data != null) {
        try {
          final String name = new String(data, ENCODING);
          return JobSchedulingStatus.valueOf(name);
        }
        catch (IllegalArgumentException e) {
          throw new IOException("Found corrupt data in ZooKeeper: " +
              Arrays.toString(data) + " is no valid job status");
        }
      }
    }
    // nothing found, yet, must be in status 'PENDING'
    return JobSchedulingStatus.PENDING;
  }
  catch (Exception e) {
    throw new IOException("Get finished state from zk fail for job " + jobID.toString(), e);
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime_2.11

@Override
public void setJobRunning(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    writeEnumToZooKeeper(jobID, JobSchedulingStatus.RUNNING);
  }
  catch (Exception e) {
    throw new IOException("Failed to set RUNNING state in ZooKeeper for job " + jobID, e);
  }
}

代码示例来源:origin: com.alibaba.blink/flink-runtime

private void writeEnumToZooKeeper(JobID jobID, JobSchedulingStatus status) throws Exception {
    final String zkPath = createZkPath(jobID);
    this.client.newNamespaceAwareEnsurePath(zkPath).ensure(client.getZookeeperClient());
    this.client.setData().forPath(zkPath, status.name().getBytes(ENCODING));
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime_2.11

@Override
public void setJobFinished(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    writeEnumToZooKeeper(jobID, JobSchedulingStatus.DONE);
  }
  catch (Exception e) {
    throw new IOException("Failed to set DONE state in ZooKeeper for job " + jobID, e);
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime

@Override
public JobSchedulingStatus getJobSchedulingStatus(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    final String zkPath = createZkPath(jobID);
    final Stat stat = client.checkExists().forPath(zkPath);
    if (stat != null) {
      // found some data, try to parse it
      final byte[] data = client.getData().forPath(zkPath);
      if (data != null) {
        try {
          final String name = new String(data, ENCODING);
          return JobSchedulingStatus.valueOf(name);
        }
        catch (IllegalArgumentException e) {
          throw new IOException("Found corrupt data in ZooKeeper: " +
              Arrays.toString(data) + " is no valid job status");
        }
      }
    }
    // nothing found, yet, must be in status 'PENDING'
    return JobSchedulingStatus.PENDING;
  }
  catch (Exception e) {
    throw new IOException("Get finished state from zk fail for job " + jobID.toString(), e);
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime_2.10

@Override
public void setJobRunning(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    writeEnumToZooKeeper(jobID, JobSchedulingStatus.RUNNING);
  }
  catch (Exception e) {
    throw new IOException("Failed to set RUNNING state in ZooKeeper for job " + jobID, e);
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime_2.11

@Override
public void clearJob(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    final String zkPath = createZkPath(jobID);
    this.client.newNamespaceAwareEnsurePath(zkPath).ensure(client.getZookeeperClient());
    this.client.delete().forPath(zkPath);
  }
  catch (Exception e) {
    throw new IOException("Failed to clear job state from ZooKeeper for job " + jobID, e);
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime_2.10

@Override
public void setJobFinished(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    writeEnumToZooKeeper(jobID, JobSchedulingStatus.DONE);
  }
  catch (Exception e) {
    throw new IOException("Failed to set DONE state in ZooKeeper for job " + jobID, e);
  }
}

代码示例来源:origin: org.apache.flink/flink-runtime

@Override
public void clearJob(JobID jobID) throws IOException {
  checkNotNull(jobID);
  try {
    final String zkPath = createZkPath(jobID);
    this.client.newNamespaceAwareEnsurePath(zkPath).ensure(client.getZookeeperClient());
    this.client.delete().forPath(zkPath);
  }
  catch (Exception e) {
    throw new IOException("Failed to clear job state from ZooKeeper for job " + jobID, e);
  }
}

相关文章