- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我需要创建一个 Apache Beam (Java) 流作业,该作业应每 60 秒启动一次(且仅启动一次)。
我通过使用GenerateSequence、Window 和Combine,使用DirectRunner 使其正常工作。
但是,当我在 Google Dataflow 上运行它时,有时它会在 60 秒窗口内触发多次。我猜这与延迟和乱序消息有关。
Pipeline pipeline = Pipeline.create(options);
pipeline
// Jenerate a tick every 15 seconds
.apply("Ticker", GenerateSequence.from(0).withRate(1, Duration.standardSeconds(15)))
// Just to check if individual ticks are being generated once every 15 second
.apply(ParDo.of(new DoFn<Long, Long>() {
@ProcessElement
public void processElement(@Element Long tick, OutputReceiver<Long> out) {
ZonedDateTime currentInstant = Instant.now().atZone(ZoneId.of("Asia/Jakarta"));
LOG.warn("-" + tick + "-" + currentInstant.toString());
out.output(word);
}
}
))
// 60 Second window
.apply("Window", Window.<Long>into(FixedWindows.of(Duration.standardSeconds(60))))
// Emit once per 60 second
.apply("Cobmine window into one", Combine.globally(Count.<Long>combineFn()).withoutDefaults())
.apply("START", ParDo.of(new DoFn<Long, ZonedDateTime>() {
@ProcessElement
public void processElement(@Element Long count, OutputReceiver<ZonedDateTime> out) {
ZonedDateTime currentInstant = Instant.now().atZone(ZoneId.of("Asia/Jakarta"));
// LOG just to check
// This log is sometimes printed more than once within 60 seconds
LOG.warn("x" + count + "-" + currentInstant.toString());
out.output(currentInstant);
}
}
));
它在大多数情况下都有效,除了随机每 5 或 10 分钟一次,我在同一分钟内看到两个输出。如何确保上面的“START”每 60 秒运行一次?谢谢。
最佳答案
简短回答:目前还不能,Beam 模型专注于事件时处理和后期数据的正确处理。
解决方法:您可以定义一个处理时间计时器,但您必须手动处理计时器和延迟数据的输出和处理,see this或this .
更多详细信息:
Beam 中的窗口和触发器通常在事件时间中定义,而不是在处理时间中定义。这样,如果您在发出窗口结果后有延迟数据,则延迟数据仍然会出现在正确的窗口中,并且可以为该窗口重新计算结果。梁模型允许您表达该逻辑,并且其大部分功能都是为此量身定制的。
这也意味着通常不需要 Beam 管道在某个特定的现实时间(例如说“根据事件本身中的数据聚合属于某个窗口的事件,然后每分钟输出该窗口”这样的说法是没有意义的。 Beam runner 聚合窗口的数据,可能会等待较晚的数据,然后在它认为正确时立即发出结果。数据准备好发出的条件由触发器指定。但这只是 - 当窗口数据准备好发出时的条件,它实际上并不强制运行器发出它。因此,运行程序可以在满足触发条件后的任何时间点发出它,并且结果将是正确的,即,如果自满足计时器条件以来有更多事件到达,则仅处理属于具体窗口的事件在那个窗口中。
事件时间窗口不能与处理时间触发一起使用,并且 Beam 中没有方便的原语(触发器/窗口)来处理存在延迟数据的处理时间。在此模型中,如果您使用仅触发一次的触发器,则会丢失最新数据,并且仍然无法定义强大的处理时间触发器。要构建这样的东西,您必须能够指定诸如现实生活中的时间点之类的内容,从该时间点开始测量处理时间,并且您将必须处理不同处理时间和可能在整个过程中发生的延迟的问题。大量 worker 机器。目前这还不是 Beam 的一部分。
Beam 社区做出了一些努力来实现此用例,例如sink triggers和 retractions这将允许您在事件时间空间中定义管道,但无需复杂的事件时间触发器。结果可以立即更新/重新计算并发出,或者可以在接收器处指定触发器,例如“我希望输出表每分钟更新一次”。结果将自动更新并重新计算最新数据,无需您的参与。但目前这些努力还远未完成,因此您目前最好的选择是使用 existing triggers 之一或使用 timers 手动处理所有内容。
关于java - 如何创建在固定时间间隔内触发一次且仅触发一次的流式Beam管道,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/57265163/
我正在使用 Assets 管道来管理我的 Grails 3.0 应用程序的前端资源。但是,似乎没有创建 CoffeeScript 文件的源映射。有什么办法可以启用它吗? 我的 build.gradle
我有一个我想要的管道: 提供一些资源, 运行一些测试, 拆资源。 我希望第 3 步中的拆卸任务运行 不管 测试是否通过或失败,在第 2 步。据我所知 runAfter如果前一个任务成功,则只运行一个任
如果我运行以下命令: Measure-Command -Expression {gci -Path C:\ -Recurse -ea SilentlyContinue | where Extensio
我知道管道是一个特殊字符,我需要使用: Scanner input = new Scanner(System.in); String line = input.next
我再次遇到同样的问题,我有我的默认处理方式,但它一直困扰着我。 有没有更好的办法? 所以基本上我有一个运行的管道,在管道内做一些事情,并想从管道内返回一个键/值对。 我希望整个管道返回一个类型为 ps
我有三个环境:dev、hml 和 qa。 在我的管道中,根据分支,阶段有一个条件来检查它是否会运行: - stage: Project_Deploy_DEV condition: eq(varia
我有 Jenkins Jenkins ver. 2.82 正在运行并想在创建新作业时使用 Pipeline 功能。但我没有看到这个列为选项。我只能在自由式项目、maven 项目、外部项目和多配置之间进
在对上一个问题 (haskell-data-hashset-from-unordered-container-performance-for-large-sets) 进行一些观察时,我偶然发现了一个奇
我正在寻找有关如何使用管道将标准输出作为其他命令的参数传递的见解。 例如,考虑这种情况: ls | grep Hello grep 的结构遵循以下模式:grep SearchTerm PathOfFi
有没有办法不因声明性管道步骤而失败,而是显示警告?目前我正在通过添加 || exit 0 来规避它到 sh 命令行的末尾,所以它总是可以正常退出。 当前示例: sh 'vendor/bin/phpcs
我们正在从旧的 Jenkins 设置迁移到所有计划都是声明性 jenkinsfile 管道的新服务器……但是,通过使用管道,我们无法再手动清除工作区。我如何设置 Jenkins 以允许 手动点播清理工
我在 Python 中阅读了有关 Pipelines 和 GridSearchCV 的以下示例: http://www.davidsbatista.net/blog/2017/04/01/docume
我有一个这样的管道脚本: node('linux'){ stage('Setup'){ echo "Build Stage" } stage('Build'){ echo
我正在使用 bitbucket 管道进行培训 这是我的 bitbucket-pipelines.yml: image: php:7.2.9 pipelines: default:
我正在编写一个程序,其中输入文件被拆分为多个文件(Shamir 的 secret 共享方案)。 这是我想象的管道: 来源:使用 Conduit.Binary.sourceFile 从输入中读取 导管:
我创建了一个管道,它有一个应该只在开发分支上执行的阶段。该阶段还需要用户输入。即使我在不同的分支上,为什么它会卡在这些步骤的用户输入上?当我提供输入时,它们会被正确跳过。 stage('Deplo
我正在尝试学习管道功能(%>%)。 当试图从这行代码转换到另一行时,它不起作用。 ---- R代码--原版----- set.seed(1014) replicate(6,sample(1:8))
在 Jenkins Pipeline 中,如何将工件从以前的构建复制到当前构建? 即使之前的构建失败,我也想这样做。 最佳答案 Stuart Rowe 还在 Pipeline Authoring Si
我正在尝试使用 执行已定义的作业构建 使用 Jenkins 管道的方法。 这是一个简单的例子: build('jenkins-test-project-build', param1 : 'some-
当我使用 where 过滤器通过管道命令排除对象时,它没有给我正确的输出。 PS C:\Users\Administrator> $proall = Get-ADComputer -filter *
我是一名优秀的程序员,十分优秀!