本文整理了Java中org.apache.flink.runtime.highavailability.zookeeper.ZooKeeperRunningJobsRegistry
类的一些代码示例,展示了ZooKeeperRunningJobsRegistry
类的具体用法。这些代码示例主要来源于Github
/Stackoverflow
/Maven
等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。ZooKeeperRunningJobsRegistry
类的具体详情如下:
包路径:org.apache.flink.runtime.highavailability.zookeeper.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);
}
}
内容来源于网络,如有侵权,请联系作者删除!