- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我想知道是否指定了流拓扑处理消息的顺序。
示例:
// read input messages
KStream<String, String> inputMessages = builder.stream("demo_input_topic_1");
inputMessages = inputMessages.peek((k, v) -> System.out.println("TECHN. NEW MESSAGE: key: " + k + ", value: " + v));
// check if message was already processed
KTable<String, Long> alreadyProcessedMessages = inputMessages.groupByKey().count();
KStream<String, String> newMessages =
inputMessages.leftJoin(alreadyProcessedMessages, (streamValue, tableValue) -> getMessageValueOrNullIfKnownMessage(streamValue, tableValue));
KStream<String, String> filteredNewMessages =
newMessages.filter((key, val) -> val != null).peek((k, v) -> System.out.println("FUNC. NEW MESSAGE: key: " + k + ", value: " + v));
// process the message
filteredNewMessages.map((key, value) -> KeyValue.pair(key, "processed message: " + value))
.peek((k, v) -> System.out.println("PROCESSED MESSAGE: key: " + k + ", value: " + v)).to("demo_output_topic_1");
使用getMessageValueOrNullIfKnownMessage(...)
:
private static String getMessageValueOrNullIfKnownMessage(String newMessageValue, Long messageCounter) {
if (messageCounter > 1) {
return null;
}
return newMessageValue;
}
因此示例中只有一个输入和一个输出主题。
输入主题在 alreadyProcessedMessages
中进行计数(从而创建本地状态)。此外,输入主题与计数表 alreadyProcessedMessages
连接,连接结果是流 newMessages
(该流中的消息值为 null
如果消息计数 > 1,否则为消息的原始值)。
然后,newMessages
的消息被过滤(null
值被过滤掉),并将结果写入输出主题。
这个最小流的作用是:它将所有消息从输入主题写入具有新 key (之前未处理过的 key )的输出主题。
在测试中,流可以工作。但我认为这并不能保证。它之所以有效,是因为消息在加入之前首先由计数节点处理。
但是该订单有保证吗?
据我在所有文档中看到的,无法保证此处理顺序。因此,如果有新消息到达,也可能会发生这种情况:
这当然会产生不同的结果(因此在这种情况下,如果具有相同 key 的消息第二次到达,它仍然会与原始值连接,因为它尚未被计数)。
那么处理顺序是在某处指定的吗?
我知道在新版本的 Kafka 中,KStream-KTable 连接是根据输入分区中消息的时间戳完成的。但这在这里没有帮助,因为拓扑使用相同的输入分区(因为它是相同的消息)。
谢谢
最佳答案
没有任何保证。即使在当前实现中,使用子节点的 List
: https://github.com/apache/kafka/blob/3.6/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java#L267-L269 -- 但是,不能保证子节点按照 DSL 中指定的顺序附加到此列表(因为中间有一个转换层,可能会以不同的顺序添加节点)。此外,实现可能随时发生变化。
我能想到的唯一解决方法(相当昂贵)可能是在重新分区主题中发送流端数据:
KStream<String, String> newMessages =
inputMessages.through(...) // note: as of 2.6.0 release, you could use `repartition()` instead of `through()`
.leftJoin(alreadyProcessedMessages, ...);
这样,KTable 将在执行连接之前更新,因为需要先读回记录。但是,由于您无法保证何时读回记录,因此在完成联接之前可能会对表进行多次更新,这可能会让您处于与以前类似的情况。 (此外,通过其他主题重新路由数据的成本有些昂贵。)
使用处理器 API,您将拥有移动控制权,因为您可以调用 context.forward(..., To.child(...))
。但是,对于这种情况,您还需要手动实现聚合和连接:
KStream routing = inputMessages.transform(...);
routing.groupByKey(...);
routing.leftJoin(...);
对于这种情况,您会在 transform()
之后获得您想要避免的重新分区主题:
KStream routing = inputMessages.transform(...);
routing.transform(...); // implement the aggregation
routing.transform(...); // implement the join
连续的 transform()
将不会触发自动重新分区。
关于apache-kafka - Kafka Streams拓扑的处理顺序是否指定?,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/61748676/
我正在我的 java 作业中使用 GUI,并且我必须指定 JCheckBox 中的其他内容。除了这个小要求,其他的我都完成了。我不太确定如何解决这个问题,我查阅了我的书并尝试在线研究 要求: 一系列复
在各种语言中(我将在这里使用 JavaScript,但我已经在 PHP 和 C++ 中以及可能在其他地方看到过它),似乎有几种构造简单 for 循环的方法。版本 1 如下: var top = doc
有没有一种方法可以使用 CSS 指定每次“小于符号”(在键盘上 M 的右侧)或“大于符号”出现在文本中时,它应该被替换为分别是“小于”或“大于”的实际词? 最佳答案 CSS 不能作用于(不能修改,即)
首先,使用 setspn 命令为用户注册服务主体名称。 setspn -a CS/dummy@abc.com dummyuser setspn -l dummyuser 给出输出为 CS/dummy@
我在指定从 SFSafariViewController 访问时遇到问题,因为它具有与 Safari 浏览器完全相同的用户代理。 我要做的是仅在 webview 内显示图片,如果在普通浏览器上查看,则
我正在尝试用 R 语言在 lavaan 中指定一个奇怪的模型。该模型如下所示: 我的规范尝试如下所示。我发现难以实现的是将观察到的变量的唯一误差固定为唯一项的两个相关性的总和。 例如,项目 y*1,2
我正在构建 API 以将我的 React 应用程序与我的后端服务连接起来,我想使用 typescript 来指定 data 的类型在我的 Axios 请求中。如何在不修改其他字段的情况下更新 Axio
如何为模型指定初始“软”值?该初始模型是解决类似查询的结果,并且该模型很可能具有正确的部分,甚至对于当前查询可能是正确的。 目前,我正在通过增量求解和 hard/soft constraints 对此
我有来自网页的以下代码 https://cwiki.apache.org/confluence/display/KAFKA/0.8.0+Producer+Example 似乎缺少的是如何配置分区数。我
有没有办法在每个查询的基础上在 Neo4jClient 中指定 Cypher 解析器的版本,如 here 所述? 谢谢! 最佳答案 如果您将 Neo4jClient 更新到最新版本(> 1.0.0.6
我有以下代码生成四个图,但它们最终被压扁(见下图)。我该如何解决这个问题? par(mfrow=c(2,2)) curve(.5*exp(-.5*x),from=0,to=10,main="f(x)"
我有一个 ColdFusion 10 服务器。我正在使用 JDBC 驱动程序连接到 db2 数据库。我偶然发现了这个笔记。这个设置在哪里?我还查看了 neo*.xml 文件,但没有看到任何 db 驱动
我想知道是否可以指定验证器的运行顺序。 目前,我编写了一个自定义验证器,检查它是否为 [a-zA-Z0-9]+ 以确保登录验证我们的规则,并编写了一个远程验证器以确保登录可用,但目前远程验证器已启动在
我的应用程序需要至少 40MB 的 RAM,因此早期的 iPhone(例如 3G、第一个 iPod touch 版本)就没有它(它们为我的应用程序提供的最大内存约为 20MB)。有没有正确的方法来禁用
我有一个保存日期(不是当前日期)的 Date 对象,我需要以某种方式指定该日期为 UTC,然后将其转换为“欧洲/巴黎”,即 +1 小时。 public static LocalDateTime toL
我想问你在 Varnish 代码中如何在没有缓存的情况下将请求传递到后端。 我知道我可以做到并且正在发挥作用: if (req.url ~ "(\?|&)(something|somethin
我目前基于模块编译程序(如主程序 foo 依赖于模块 bar )如下: gfortran -c bar.f90 gfortran -o foo.exe foo.f90 bar.o 这在 foo.f90
我正在尝试创建一个依赖于另一个 meteor 包的新 meteor 包。当我尝试 meteor add mypackage 时,出现以下错误。为什么 Meteor 不添加 mypackage 并引入它
我正在制作执行器/ react 器,同时发现这是一个终生的问题。它与 async/Future 无关,可以在没有 async 糖的情况下进行复制。 use std::future::Future; s
我在 cassandra 中有一个表,其数据类型为时间戳。我正在使用 cqlsh 从数据库中获取数据,并希望更改我的时间戳列输出的输出格式。我研究了一下,发现我可以通过更改以下文件来更改时间戳输出格式
我是一名优秀的程序员,十分优秀!