- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我目前正在开发一个 kafka java 项目。我是新手,我发现很难理解与 Kafka 生产者/消费者设计相关的一些基本概念。
假设,我有一个具有单个分区的主题,并且有一个生产者写入该主题,还有一个消费者从该主题消费。如果我部署同一应用程序的多个实例,每个实例都将运行它自己的使用者。这样的话,由于所有的消费者都属于同一个groupId,那么消息会均匀地分布在运行在多个实例上的消费者之间吗?
如何从应用程序中定期检查消费者是否存活?
请对上述疑问进行澄清。如果我的任何/所有假设/理解是错误的,请纠正我。我知道我没有分享任何代码示例,因为这些是概念性问题。如果需要,我可以分享代码片段。
最佳答案
您说具有单个分区的主题意味着它无法将消息分发到多个分区。你将失去 Kafka 的一大优势。您必须将分区增加一个以上。如果您部署同一应用程序的多个实例,则它无助于分发,因为正如您提到的,消息将发布到一个分区,并且只有一个实例只会分配给该分区,其他实例将处于空闲状态。
<您可以使用 AdminClient Kafka API 检查您的消费者是否有任何滞后
Properties props = new Properties();
props.setProperty("bootstrap.servers", "localhost:9091");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
AdminClient client = org.apache.kafka.clients.admin.AdminClient.create(props);
ListConsumerGroupOffsetsResult offsets = client.listConsumerGroupOffsets("consumerId");
Map<TopicPartition, OffsetAndMetadata> tt = offsets.partitionsToOffsetAndMetadata().get();
ListConsumerGroupOffsetsResult offsets = client.listConsumerGroupOffsets(consumerId);
Map<TopicPartition, OffsetAndMetadata> tt = offsets.partitionsToOffsetAndMetadata().get();
for (Entry<TopicPartition, OffsetAndMetadata> entry : tt.entrySet()) {TopicPartition tp = entry.getKey();
OffsetAndMetadata op = entry.getValue();
Collections.singletonList(tp);
consumer.assign(Collections.singletonList(tp));
consumer.seekToEnd(Collections.singletonList(tp));
System.out.println(consumerId + "," + tp.partition() + "," + consumer.position(tp) + ","
+ op.offset() + "," + (consumer.position(tp) - op.offset()));
}
您没有说明您的部署位置,但如果您在 mesos 中使用 marathon 进行部署,它将自动重新启动。您可以手动重新启动,如果您使用与之前相同的组 ID,您的应用程序将开始使用它留下的位置。
关于java - Kafka - 如何检查消费者是否还活着,如果不是,如何将消费者恢复到运行状态?,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/58353237/
我正在通读 Windows Phone 7.5 Unleashed,有很多代码看起来像这样(在页面的代码隐藏中): bool loaded; protected override void OnNav
在cgi服务器中,我这样返回 print ('Status: 201 Created') print ('Content-Type: text/html') print ('Location: htt
我正在查看 esh(easy shell)的实现,无法理解在这种情况下什么是 22 和 9 信号。理想情况下,有一个更具描述性的常量,但我找不到列表。 最佳答案 信号列表及其编号(包括您看到的这两个)
我的Oozie Hive Action 永远处于运行模式。 oozie.log文件中没有显示错误。
我正在编写一个使用 RFCOMM 通过蓝牙连接到设备的 Android 应用程序。我使用 BluetoothChat 示例作为建立连接的基础,大部分时间一切正常。 但是,有时由于出现套接字已打开的消息
我有一个云调度程序作业,它应该每小时访问我的 API 以更新一些价格。这些作业大约需要 80 秒才能运行。 这是它的作用: POST https://www.example.com/api/jobs/
我正在 Tomcat 上访问一个简单的 JSP 页面: 但是当我使用 curl 测试此页面时,我得到了 200 响应代码而不是预期的 202: $ curl -i "http://localhos
有时 JAR-RS 客户端会发送错误的语法请求正文。服务器应响应 HTTP status 400 (Bad Request) , 但它以 HTTP status 500 (Internal Serve
我正在尝试通过 response.send() 发送一个整数,但我不断收到此错误 express deprecated res.send(status): Use res.sendStatus(sta
我已经用 Excel 和 Java 做过很多次了……这次我需要用 Stata 来做,因为保存变量更方便'labels .如何将 dataset_1 重组为下面的 dataset_2? 我需要转换以下
我正在创建一个应用程序,其中的对象具有状态查找功能。为了提供一些上下文,让我们使用以下示例。 帮助台应用程序,其中创建作业并通过以下工作流程移动: 新 - 工作已创建但未分配 进行中 - 分配给工作人
我想在 Keras 中运行 LSTM 并获得输出和状态。在 TF 中有这样的事情 with tf.variable_scope("RNN"): for time_step in range
有谁知道 Scala-GWT 的当前状态 项目? 那里的主要作者 Grzegorz Kossakowski 似乎退出了这个项目,在 Spring 中从事 scalac 的工作。 但是,在 interv
我正在尝试编写一个 super 简单的 applescript 来启动 OneDrive App , 或确保打开,当机器的电源设置为插入时,将退出,或确保关闭,当电源设置为电池时。 我无法找到如何访问
目前我正在做这样的事情 link.on('click', function () { if (link.attr('href') !== $route.current.originalPath
是否可以仅通过查看用户代理来检测浏览器上是否启用/禁用 Javascript。 如果是,我应该寻找什么。如果否,检测用户浏览器是否启用/禁用 JavaScript 的最佳方法是什么 最佳答案 不,没有
Spring 和 OSGi 目前的开发状况如何? 最近好像有点安静了。 文档的最新版本 ( http://docs.spring.io/osgi/ ) 来自 2009 年。 我看到一些声明 Sprin
我正在从主函数为此类创建一个线程,但即使使用 Thread.currentThread().interrupt() 中断它,输出仍然包含“Still Here”行。 public class Writ
为了满足并发要求,我想知道如何在 Godog 中的多个步骤之间传递参数或状态。 func FeatureContext(s *godog.Suite) { // This step is ca
我有一个UIButton子类,它不使用UIImage背景,仅使用背景色。我注意到的一件事是,当您设置按钮的背景图像时,有一个默认的突出显示状态,当按下按钮时,该按钮会稍微变暗。 这是我当前的代码。
我是一名优秀的程序员,十分优秀!