- 使用 Spring Initializr 创建 Spring Boot 应用程序
- 在Spring Boot中配置Cassandra
- 在 Spring Boot 上配置 Tomcat 连接池
- 将Camel消息路由到嵌入WildFly的Artemis上
本文整理了Java中org.apache.samza.zk.ZkBarrierForVersionUpgrade
类的一些代码示例,展示了ZkBarrierForVersionUpgrade
类的具体用法。这些代码示例主要来源于Github
/Stackoverflow
/Maven
等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。ZkBarrierForVersionUpgrade
类的具体详情如下:
包路径:org.apache.samza.zk.ZkBarrierForVersionUpgrade
类名称:ZkBarrierForVersionUpgrade
[英]ZkBarrierForVersionUpgrade is an implementation of distributed barrier, which guarantees that the expected barrier size and barrier participants match before marking the barrier as complete. It also allows the caller to expire the barrier. This implementation is specifically tailored towards barrier support during jobmodel version upgrades. The participant responsible for the version upgrade starts the barrier by invoking #create(String,List). Each participant in the list, then, joins the new barrier. When all listed participants #join(String,String)the barrier, the creator marks the barrier as org.apache.samza.zk.ZkBarrierForVersionUpgrade.State#DONEwhich signals the end of barrier. The creator of the barrier can expire the barrier by invoking #expire(String). This will mark the barrier with value org.apache.samza.zk.ZkBarrierForVersionUpgrade.State#TIMED_OUT and indicates to everyone that it is no longer valid. Describes the lifecycle of a barrier.
When expected participants join
Leader ---< NEW ---------------------------------------- < DONE
| barrier within barrierTimeOut.
|
|
|
|
|
| When expected participants doesn't
| ----------------------------------------- < TIMED_OUT
join barrier within barrierTimeOut.
The caller can listen to events associated with the barrier by registering a ZkBarrierListener. Zk Tree Reference: /barrierRoot/ | |- barrier_{version1}/ | |- barrier_state/ | | ([NEW|DONE|TIMED_OUT]) | |- barrier_participants/ | | |- {id1} | | |- {id2} | | |- ...
[中]ZkBarrierForVersionUpgrade是分布式屏障的一种实现,它保证在标记屏障完成之前,预期的屏障大小和屏障参与者匹配。它还允许调用者终止障碍。此实现专门针对jobmodel版本升级期间的屏障支持而定制。负责版本升级的参与者通过调用#create(String,List)启动屏障。然后,列表中的每个参与者都加入了新的障碍。当所有列出的参与者#加入(String,String)障碍时,创建者将障碍标记为org。阿帕奇。萨姆萨。zk。ZkBarrierForVersionUpgrade。状态#完成,表示屏障结束。屏障的创建者可以通过调用#expire(字符串)使屏障过期。这将用价值组织标记障碍。阿帕奇。萨姆萨。zk。ZkBarrierForVersionUpgrade。状态#超时#并向所有人表明它不再有效。描述屏障的生命周期
When expected participants join
Leader ---< NEW ---------------------------------------- < DONE
| barrier within barrierTimeOut.
|
|
|
|
|
| When expected participants doesn't
| ----------------------------------------- < TIMED_OUT
join barrier within barrierTimeOut.
调用者可以通过注册ZKBrierListener来监听与屏障相关的事件。Zk树引用:/barrieroot/| | |-barrier{version1}/| |-barrier | state/| |([NEW | DONE | TIMED |])|-barrier |参与者| | | |-{id1 | | |-{id2 | |-。。。
代码示例来源:origin: apache/samza
@Override
public int compare(String o1, String o2) {
// barrier's name format is barrier_<num>
return ZkBarrierForVersionUpgrade.getVersion(o1) - ZkBarrierForVersionUpgrade.getVersion(o2);
}
});
代码示例来源:origin: apache/samza
ZkBarrierForVersionUpgrade processor1Barrier = new ZkBarrierForVersionUpgrade(barrierId, zkUtils, listener, debounceTimer);
ZkBarrierForVersionUpgrade processor2Barrier = new ZkBarrierForVersionUpgrade(barrierId, zkUtils1, listener, debounceTimer);
processor1Barrier.create(BARRIER_VERSION, processors);
processor1Barrier.join(BARRIER_VERSION, "p1");
processor2Barrier.join(BARRIER_VERSION, "p2");
processor1Barrier.expire(BARRIER_VERSION);
processor1Barrier.join(BARRIER_VERSION, "p3");
代码示例来源:origin: apache/samza
@Test
public void testZkBarrierForVersionUpgrade() {
String barrierId = String.format("%s/%s", zkUtils1.getKeyBuilder().getRootPath(), RandomStringUtils.randomAlphabetic(4));
List<String> processors = ImmutableList.of("p1", "p2");
CountDownLatch latch = new CountDownLatch(2);
TestZkBarrierListener listener = new TestZkBarrierListener(latch, State.DONE);
ZkBarrierForVersionUpgrade processor1Barrier = new ZkBarrierForVersionUpgrade(barrierId, zkUtils, listener, debounceTimer);
ZkBarrierForVersionUpgrade processor2Barrier = new ZkBarrierForVersionUpgrade(barrierId, zkUtils1, listener, debounceTimer);
processor1Barrier.create(BARRIER_VERSION, processors);
State barrierState = zkUtils.getZkClient().readData(barrierId + "/barrier_1/barrier_state");
assertEquals(State.NEW, barrierState);
processor1Barrier.join(BARRIER_VERSION, "p1");
processor2Barrier.join(BARRIER_VERSION, "p2");
boolean result = false;
try {
result = latch.await(10000, TimeUnit.MILLISECONDS);
} catch (InterruptedException e) {
e.printStackTrace();
}
assertTrue("Barrier failed to complete within test timeout.", result);
List<String> children = zkUtils.getZkClient().getChildren(barrierId + "/barrier_1/barrier_participants");
barrierState = zkUtils.getZkClient().readData(barrierId + "/barrier_1/barrier_state");
assertEquals(State.DONE, barrierState);
assertNotNull(children);
assertEquals("Unexpected barrier state. Didn't find two processors.", 2, children.size());
assertEquals("Unexpected barrier state. Didn't find the expected members.", processors, children);
}
代码示例来源:origin: apache/samza
@Test
public void testCleanUpZkBarrierVersion() {
String root = zkUtils.getKeyBuilder().getJobModelVersionBarrierPrefix();
zkUtils.getZkClient().createPersistent(root, true);
ZkBarrierForVersionUpgrade barrier = new ZkBarrierForVersionUpgrade(root, zkUtils, null, null);
for (int i = 200; i < 210; i++) {
barrier.create(String.valueOf(i), new ArrayList<>(Arrays.asList(i + "a", i + "b", i + "c")));
}
zkUtils.deleteOldBarrierVersions(5);
List<String> zNodeIds = zkUtils.getZkClient().getChildren(root);
Collections.sort(zNodeIds);
Assert.assertEquals(Arrays.asList("barrier_205", "barrier_206", "barrier_207", "barrier_208", "barrier_209"),
zNodeIds);
}
代码示例来源:origin: org.apache.samza/samza-core
@Override
public void onBarrierCreated(String version) {
// Start the timer for rebalancing
startTime = System.nanoTime();
metrics.barrierCreation.inc();
if (leaderElector.amILeader()) {
debounceTimer.scheduleAfterDebounceTime(barrierAction, (new ZkConfig(config)).getZkBarrierTimeoutMs(), () -> barrier.expire(version));
}
}
代码示例来源:origin: org.apache.samza/samza-core_2.11
ZkJobCoordinator(Config config, MetricsRegistry metricsRegistry, ZkUtils zkUtils) {
this.config = config;
this.metrics = new ZkJobCoordinatorMetrics(metricsRegistry);
this.processorId = createProcessorId(config);
this.zkUtils = zkUtils;
// setup a listener for a session state change
// we are mostly interested in "session closed" and "new session created" events
zkUtils.getZkClient().subscribeStateChanges(new ZkSessionStateChangedListener());
leaderElector = new ZkLeaderElector(processorId, zkUtils);
leaderElector.setLeaderElectorListener(new LeaderElectorListenerImpl());
this.debounceTimeMs = new JobConfig(config).getDebounceTimeMs();
this.reporters = MetricsReporterLoader.getMetricsReporters(new MetricsConfig(config), processorId);
debounceTimer = new ScheduleAfterDebounceTime(processorId);
debounceTimer.setScheduledTaskCallback(throwable -> {
LOG.error("Received exception in debounce timer! Stopping the job coordinator", throwable);
stop();
});
this.barrier = new ZkBarrierForVersionUpgrade(zkUtils.getKeyBuilder().getJobModelVersionBarrierPrefix(), zkUtils, new ZkBarrierListenerImpl(), debounceTimer);
systemAdmins = new SystemAdmins(config);
streamMetadataCache = new StreamMetadataCache(systemAdmins, METADATA_CACHE_TTL_MS, SystemClock.instance());
}
代码示例来源:origin: org.apache.samza/samza-core_2.12
/**
* Invoked when there is a change to the JobModelVersion z-node. It signifies that a new JobModel version is available.
*/
@Override
public void doHandleDataChange(String dataPath, Object data) {
debounceTimer.scheduleAfterDebounceTime(JOB_MODEL_VERSION_CHANGE, 0, () -> {
String jobModelVersion = (String) data;
LOG.info("Got a notification for new JobModel version. Path = {} Version = {}", dataPath, data);
newJobModel = zkUtils.getJobModel(jobModelVersion);
LOG.info("pid=" + processorId + ": new JobModel is available. Version =" + jobModelVersion + "; JobModel = " + newJobModel);
if (!newJobModel.getContainers().containsKey(processorId)) {
LOG.info("New JobModel does not contain pid={}. Stopping this processor. New JobModel: {}",
processorId, newJobModel);
stop();
} else {
// stop current work
if (coordinatorListener != null) {
coordinatorListener.onJobModelExpired();
}
// update ZK and wait for all the processors to get this new version
barrier.join(jobModelVersion, processorId);
}
});
}
代码示例来源:origin: org.apache.samza/samza-core
barrier.create(nextJMVersion, currentProcessorIds);
代码示例来源:origin: apache/samza
@Test
public void testShouldDiscardBarrierUpdateEventsAfterABarrierIsMarkedAsDone() {
String barrierId = String.format("%s/%s", zkUtils1.getKeyBuilder().getRootPath(), RandomStringUtils.randomAlphabetic(4));
List<String> processors = ImmutableList.of("p1", "p2");
CountDownLatch latch = new CountDownLatch(2);
TestZkBarrierListener listener = new TestZkBarrierListener(latch, State.DONE);
ZkBarrierForVersionUpgrade processor1Barrier = new ZkBarrierForVersionUpgrade(barrierId, zkUtils, listener, debounceTimer);
ZkBarrierForVersionUpgrade processor2Barrier = new ZkBarrierForVersionUpgrade(barrierId, zkUtils1, listener, debounceTimer);
processor1Barrier.create(BARRIER_VERSION, processors);
processor1Barrier.join(BARRIER_VERSION, "p1");
processor2Barrier.join(BARRIER_VERSION, "p2");
boolean result = false;
try {
result = latch.await(10000, TimeUnit.MILLISECONDS);
} catch (InterruptedException e) {
e.printStackTrace();
}
assertTrue("Barrier Timeout test failed to complete within test timeout.", result);
processor1Barrier.expire(BARRIER_VERSION);
State barrierState = zkUtils.getZkClient().readData(barrierId + "/barrier_1/barrier_state");
assertEquals(State.DONE, barrierState);
}
代码示例来源:origin: apache/samza
@Override
public void onBarrierCreated(String version) {
// Start the timer for rebalancing
startTime = System.nanoTime();
metrics.barrierCreation.inc();
if (leaderElector.amILeader()) {
debounceTimer.scheduleAfterDebounceTime(barrierAction, (new ZkConfig(config)).getZkBarrierTimeoutMs(), () -> barrier.expire(version));
}
}
代码示例来源:origin: org.apache.samza/samza-core_2.12
ZkJobCoordinator(Config config, MetricsRegistry metricsRegistry, ZkUtils zkUtils) {
this.config = config;
this.metrics = new ZkJobCoordinatorMetrics(metricsRegistry);
this.processorId = createProcessorId(config);
this.zkUtils = zkUtils;
// setup a listener for a session state change
// we are mostly interested in "session closed" and "new session created" events
zkUtils.getZkClient().subscribeStateChanges(new ZkSessionStateChangedListener());
leaderElector = new ZkLeaderElector(processorId, zkUtils);
leaderElector.setLeaderElectorListener(new LeaderElectorListenerImpl());
this.debounceTimeMs = new JobConfig(config).getDebounceTimeMs();
this.reporters = MetricsReporterLoader.getMetricsReporters(new MetricsConfig(config), processorId);
debounceTimer = new ScheduleAfterDebounceTime(processorId);
debounceTimer.setScheduledTaskCallback(throwable -> {
LOG.error("Received exception in debounce timer! Stopping the job coordinator", throwable);
stop();
});
this.barrier = new ZkBarrierForVersionUpgrade(zkUtils.getKeyBuilder().getJobModelVersionBarrierPrefix(), zkUtils, new ZkBarrierListenerImpl(), debounceTimer);
systemAdmins = new SystemAdmins(config);
streamMetadataCache = new StreamMetadataCache(systemAdmins, METADATA_CACHE_TTL_MS, SystemClock.instance());
}
代码示例来源:origin: apache/samza
/**
* Invoked when there is a change to the JobModelVersion z-node. It signifies that a new JobModel version is available.
*/
@Override
public void doHandleDataChange(String dataPath, Object data) {
debounceTimer.scheduleAfterDebounceTime(JOB_MODEL_VERSION_CHANGE, 0, () -> {
String jobModelVersion = (String) data;
LOG.info("Got a notification for new JobModel version. Path = {} Version = {}", dataPath, data);
newJobModel = zkUtils.getJobModel(jobModelVersion);
LOG.info("pid=" + processorId + ": new JobModel is available. Version =" + jobModelVersion + "; JobModel = " + newJobModel);
if (!newJobModel.getContainers().containsKey(processorId)) {
LOG.info("New JobModel does not contain pid={}. Stopping this processor. New JobModel: {}",
processorId, newJobModel);
stop();
} else {
// stop current work
if (coordinatorListener != null) {
coordinatorListener.onJobModelExpired();
}
// update ZK and wait for all the processors to get this new version
barrier.join(jobModelVersion, processorId);
}
});
}
代码示例来源:origin: org.apache.samza/samza-core_2.10
barrier.create(nextJMVersion, currentProcessorIds);
代码示例来源:origin: apache/samza
@Test
public void testZkBarrierForVersionUpgradeWithTimeOut() {
String barrierId = String.format("%s/%s", zkUtils1.getKeyBuilder().getRootPath(), RandomStringUtils.randomAlphabetic(4));
List<String> processors = ImmutableList.of("p1", "p2", "p3");
CountDownLatch latch = new CountDownLatch(2);
TestZkBarrierListener listener = new TestZkBarrierListener(latch, State.TIMED_OUT);
ZkBarrierForVersionUpgrade processor1Barrier = new ZkBarrierForVersionUpgrade(barrierId, zkUtils, listener, debounceTimer);
ZkBarrierForVersionUpgrade processor2Barrier = new ZkBarrierForVersionUpgrade(barrierId, zkUtils1, listener, debounceTimer);
processor1Barrier.create(BARRIER_VERSION, processors);
processor1Barrier.join(BARRIER_VERSION, "p1");
processor2Barrier.join(BARRIER_VERSION, "p2");
processor1Barrier.expire(BARRIER_VERSION);
boolean result = false;
try {
result = latch.await(10000, TimeUnit.MILLISECONDS);
} catch (InterruptedException e) {
e.printStackTrace();
}
assertTrue("Barrier Timeout test failed to complete within test timeout.", result);
List<String> children = zkUtils.getZkClient().getChildren(barrierId + "/barrier_1/barrier_participants");
State barrierState = zkUtils.getZkClient().readData(barrierId + "/barrier_1/barrier_state");
assertEquals(State.TIMED_OUT, barrierState);
assertNotNull(children);
assertEquals("Unexpected barrier state. Didn't find two processors.", 2, children.size());
assertEquals("Unexpected barrier state. Didn't find the expected members.", ImmutableList.of("p1", "p2"), children);
}
代码示例来源:origin: org.apache.samza/samza-core_2.12
@Override
public void onBarrierCreated(String version) {
// Start the timer for rebalancing
startTime = System.nanoTime();
metrics.barrierCreation.inc();
if (leaderElector.amILeader()) {
debounceTimer.scheduleAfterDebounceTime(barrierAction, (new ZkConfig(config)).getZkBarrierTimeoutMs(), () -> barrier.expire(version));
}
}
代码示例来源:origin: org.apache.samza/samza-core_2.10
ZkJobCoordinator(Config config, MetricsRegistry metricsRegistry, ZkUtils zkUtils) {
this.config = config;
this.metrics = new ZkJobCoordinatorMetrics(metricsRegistry);
this.processorId = createProcessorId(config);
this.zkUtils = zkUtils;
// setup a listener for a session state change
// we are mostly interested in "session closed" and "new session created" events
zkUtils.getZkClient().subscribeStateChanges(new ZkSessionStateChangedListener());
leaderElector = new ZkLeaderElector(processorId, zkUtils);
leaderElector.setLeaderElectorListener(new LeaderElectorListenerImpl());
this.debounceTimeMs = new JobConfig(config).getDebounceTimeMs();
this.reporters = MetricsReporterLoader.getMetricsReporters(new MetricsConfig(config), processorId);
debounceTimer = new ScheduleAfterDebounceTime(processorId);
debounceTimer.setScheduledTaskCallback(throwable -> {
LOG.error("Received exception in debounce timer! Stopping the job coordinator", throwable);
stop();
});
this.barrier = new ZkBarrierForVersionUpgrade(zkUtils.getKeyBuilder().getJobModelVersionBarrierPrefix(), zkUtils, new ZkBarrierListenerImpl(), debounceTimer);
systemAdmins = new SystemAdmins(config);
streamMetadataCache = new StreamMetadataCache(systemAdmins, METADATA_CACHE_TTL_MS, SystemClock.instance());
}
代码示例来源:origin: org.apache.samza/samza-core_2.11
/**
* Invoked when there is a change to the JobModelVersion z-node. It signifies that a new JobModel version is available.
*/
@Override
public void doHandleDataChange(String dataPath, Object data) {
debounceTimer.scheduleAfterDebounceTime(JOB_MODEL_VERSION_CHANGE, 0, () -> {
String jobModelVersion = (String) data;
LOG.info("Got a notification for new JobModel version. Path = {} Version = {}", dataPath, data);
newJobModel = zkUtils.getJobModel(jobModelVersion);
LOG.info("pid=" + processorId + ": new JobModel is available. Version =" + jobModelVersion + "; JobModel = " + newJobModel);
if (!newJobModel.getContainers().containsKey(processorId)) {
LOG.info("New JobModel does not contain pid={}. Stopping this processor. New JobModel: {}",
processorId, newJobModel);
stop();
} else {
// stop current work
if (coordinatorListener != null) {
coordinatorListener.onJobModelExpired();
}
// update ZK and wait for all the processors to get this new version
barrier.join(jobModelVersion, processorId);
}
});
}
代码示例来源:origin: org.apache.samza/samza-core_2.12
barrier.create(nextJMVersion, currentProcessorIds);
代码示例来源:origin: org.apache.samza/samza-core_2.11
@Override
public int compare(String o1, String o2) {
// barrier's name format is barrier_<num>
return ZkBarrierForVersionUpgrade.getVersion(o1) - ZkBarrierForVersionUpgrade.getVersion(o2);
}
});
代码示例来源:origin: org.apache.samza/samza-core_2.10
@Override
public void onBarrierCreated(String version) {
// Start the timer for rebalancing
startTime = System.nanoTime();
metrics.barrierCreation.inc();
if (leaderElector.amILeader()) {
debounceTimer.scheduleAfterDebounceTime(barrierAction, (new ZkConfig(config)).getZkBarrierTimeoutMs(), () -> barrier.expire(version));
}
}
我在网上搜索但没有找到任何合适的文章解释如何使用 javascript 使用 WCF 服务,尤其是 WebScriptEndpoint。 任何人都可以对此给出任何指导吗? 谢谢 最佳答案 这是一篇关于
我正在编写一个将运行 Linux 命令的 C 程序,例如: cat/etc/passwd | grep 列表 |剪切-c 1-5 我没有任何结果 *这里 parent 等待第一个 child (chi
所以我正在尝试处理文件上传,然后将该文件作为二进制文件存储到数据库中。在我存储它之后,我尝试在给定的 URL 上提供文件。我似乎找不到适合这里的方法。我需要使用数据库,因为我使用 Google 应用引
我正在尝试制作一个宏,将下面的公式添加到单元格中,然后将其拖到整个列中并在 H 列中复制相同的公式 我想在 F 和 H 列中输入公式的数据 Range("F1").formula = "=IF(ISE
问题类似于this one ,但我想使用 OperatorPrecedenceParser 解析带有函数应用程序的表达式在 FParsec . 这是我的 AST: type Expression =
我想通过使用 sequelize 和 node.js 将这个查询更改为代码取决于在哪里 select COUNT(gender) as genderCount from customers where
我正在使用GNU bash,版本5.0.3(1)-发行版(x86_64-pc-linux-gnu),我想知道为什么简单的赋值语句会出现语法错误: #/bin/bash var1=/tmp
这里,为什么我的代码在 IE 中不起作用。我的代码适用于所有浏览器。没有问题。但是当我在 IE 上运行我的项目时,它发现错误。 而且我的 jquery 类和 insertadjacentHTMl 也不
我正在尝试更改标签的innerHTML。我无权访问该表单,因此无法编辑 HTML。标签具有的唯一标识符是“for”属性。 这是输入和标签的结构:
我有一个页面,我可以在其中返回用户帖子,可以使用一些 jquery 代码对这些帖子进行即时评论,在发布新评论后,我在帖子下插入新评论以及删除 按钮。问题是 Delete 按钮在新插入的元素上不起作用,
我有一个大约有 20 列的“管道分隔”文件。我只想使用 sha1sum 散列第一列,它是一个数字,如帐号,并按原样返回其余列。 使用 awk 或 sed 执行此操作的最佳方法是什么? Accounti
我需要将以下内容插入到我的表中...我的用户表有五列 id、用户名、密码、名称、条目。 (我还没有提交任何东西到条目中,我稍后会使用 php 来做)但由于某种原因我不断收到这个错误:#1054 - U
所以我试图有一个输入字段,我可以在其中输入任何字符,但然后将输入的值小写,删除任何非字母数字字符,留下“。”而不是空格。 例如,如果我输入: 地球的 70% 是水,-!*#$^^ & 30% 土地 输
我正在尝试做一些我认为非常简单的事情,但出于某种原因我没有得到想要的结果?我是 javascript 的新手,但对 java 有经验,所以我相信我没有使用某种正确的规则。 这是一个获取输入值、检查选择
我想使用 angularjs 从 mysql 数据库加载数据。 这就是应用程序的工作原理;用户登录,他们的用户名存储在 cookie 中。该用户名显示在主页上 我想获取这个值并通过 angularjs
我正在使用 autoLayout,我想在 UITableViewCell 上放置一个 UIlabel,它应该始终位于单元格的右侧和右侧的中心。 这就是我想要实现的目标 所以在这里你可以看到我正在谈论的
我需要与 MySql 等效的 elasticsearch 查询。我的 sql 查询: SELECT DISTINCT t.product_id AS id FROM tbl_sup_price t
我正在实现代码以使用 JSON。 func setup() { if let flickrURL = NSURL(string: "https://api.flickr.com/
我尝试使用for循环声明变量,然后测试cols和rols是否相同。如果是,它将运行递归函数。但是,我在 javascript 中执行 do 时遇到问题。有人可以帮忙吗? 现在,在比较 col.1 和
我举了一个我正在处理的问题的简短示例。 HTML代码: 1 2 3 CSS 代码: .BB a:hover{ color: #000; } .BB > li:after {
我是一名优秀的程序员,十分优秀!