【问题标题】:Closing a Kafka connection关闭 Kafka 连接
【发布时间】:2014-11-06 10:08:27
【问题描述】:

我有一个应用程序应该向 Kafka 发送有限数量的消息然后退出。出于某种原因,即使我关闭了生产者,Kafka 连接也会保持正常。我的实现(在 Scala 中)或多或少

object Kafka {
  private val props = new Properties()

  props.put("compression.codec", DefaultCompressionCodec.codec.toString)
  props.put("producer.type", "sync")
  props.put("metadata.broker.list", "localhost:9092")
  props.put("batch.num.messages", "200")
  props.put("message.send.max.retries", "3")
  props.put("request.required.acks", "-1")
  props.put("client.id", "myclient")

  private val producer = new Producer[Array[Byte], Array[Byte]](new ProducerConfig(props))
  private def encode(msg: Message) = new KeyedMessage("topic", msg.id.getBytes, write(msg).getBytes)

  def send(msg: Message) = Try(producer.send(encode(msg)))
  def close() = producer.close()
}

这里Message是一个简单的case类,我如何将它转换为字节数组并不重要。

消息确实到达了,但是当我最终调用Kafka.close() 时,应用程序没有退出,并且连接似乎没有被释放。

有没有办法明确要求 Kafka 终止连接?

【问题讨论】:

    标签: scala apache-kafka resource-leak


    【解决方案1】:

    def close() = producer.close()

    这将创建一个名为“close”的函数,该函数调用 producer.close() 我没有看到任何证据表明您的代码实际上关闭了生产者。

    您只需调用: 生产者关闭

    【讨论】:

    • 很抱歉,如果不清楚。正如我所说,我完成后调用Kafka.close(),然后调用producer.close()。但连接仍然保持打开状态
    • 您在哪里看到打开的连接?还有,哪个版本的 Kafka?
    • 版本为 Kafka 0.8.1.1。至于卡夫卡保持开放,我只是看到程序没有终止。你知道如何检查是否有任何与 Kafka 相关的东西仍然开放吗?也许使用 jvisualvm 来查看线程
    • 抱歉,原来不是 Kafka 保持打开状态,而是另一个连接
    猜你喜欢
    • 2017-02-12
    • 2019-07-27
    • 1970-01-01
    • 2023-03-08
    • 2022-06-27
    • 1970-01-01
    • 1970-01-01
    • 2020-08-11
    • 1970-01-01
    相关资源
    最近更新 更多