【发布时间】: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