【问题标题】:Custom event handling for KafkaAdminClientKafkaAdminClient 的自定义事件处理
【发布时间】:2021-06-24 01:19:13
【问题描述】:

我的目标是在经纪人倒闭时做一些事情,但无法做到。

代码很简单:

val properties = new Properties()
properties.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")
val client = AdminClient.create(properties)
//Suppose that the App just runs from here without consuming/producing

它启动了,然后我手动关闭了 kafka。

Logs arrives:
2021-06-23T13:51:16,681+02:00 WARN  [kafka-admin-client-thread | adminclient-1] org.apache.kafka.clients.NetworkClient: [AdminClient clientId=adminclient-1] Connection to node -1 (localhost/127.0.0.1:9092) could not be established. Broker may not be available.

如何处理?基本上我只想在代理关闭时调用自定义方法。

  • 没有什么我可以“抓住”

  • 甚至在 AdminClient/KafkaAdminClient 中都找不到 evenListener(或者我只是看错了地方)

编辑:当然我也想在代理恢复正常时调用我的自定义代码

【问题讨论】:

  • @RanLupovich “当经纪人倒闭时,你想做什么?”将有关此事件的注释发送到外部系统。

标签: scala apache-kafka


【解决方案1】:

您不能向未运行的服务器发出命令...您需要在不使用 Kafka 相关工具的代理服务器本身上运行 Kafka Java 进程检查(例如jps 或systemctl 或检查一些/var/run/kafka.pid)

【讨论】:

  • “您不能向未运行的服务器发出命令” 我想在代理关闭时向其发送命令的服务器仍在运行。与卡夫卡无关。 AdminClient 也仍在运行,否则它不会记录“brokerdown”警告。
  • 您说“然后 [您] 手动关闭 kafka”,这意味着 AdminClient(和所有其他客户端)将断开连接,并且您无法捕获或观看任何内容,因为所有 Kafka 客户端进程都会失败。正如所回答的,您需要对代理本身进行外部进程监控。如果你真的想要即时反馈,你也可以在 /brokers/ids 上设置一个 ZNode 手表
  • "因为所有 Kafka 客户端进程都会失败" 根本没有,它们仍在运行,同时不断记录代理停机问题。
  • 我想这取决于你使用什么框架,但我记得默认情况下,当Broker may not be available 时,一个普通的kafka-clients 应用程序将无法完全启动
猜你喜欢
  • 2018-05-28
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多