【问题标题】:What happens if I don't close the kafka producer如果我不关闭 kafka 生产者会发生什么
【发布时间】:2017-08-20 18:55:22
【问题描述】:

我正在处理 xml,我需要为每条记录发送一条消息,当我收到最后一条记录时,我关闭了 kafka 生产者,这里的问题是 kafka 生产者的发送方法是异步的,因此,有时当我关闭它的生产者java.lang.IllegalStateException: Cannot send after the producer is closed. 我在某处读过我可以让生产者保持打开状态。我的问题是:这意味着什么,或者是否有更好的解决方案。

---编辑---

<list>
  <element attr1="" att2="" attr3=""/>
  <element attr1="" att2="" attr3=""/>
  <element attr1="" att2="" attr3=""/>
  <element attr1="" att2="" attr3=""/>
  <element attr1="" att2="" attr3=""/>
  <element attr1="" att2="" attr3=""/>
  <element attr1="" att2="" attr3=""/>
  <element attr1="" att2="" attr3=""/>
...
</list>

想象以下场景:

  • 我们读取标签并创建 kafka 生产者
  • 我们读取每个元素的属性,生成一个 json 对象并使用 send 方法将其发送到 kafka。 - 当我们读取元素时,我们在生产者中调用 close 方法

问题元素的数量可能是 80k 因此,有时当我们调用 disconnect 方法时,它会继续以异步方式发送消息。所以我们需要先调用flush方法但是它会影响性能

【问题讨论】:

    标签: java apache-kafka kafka-producer-api


    【解决方案1】:

    您应该在致电Producer.close() 之前致电Producer.flush()。这是一个阻塞调用,在发送所有记录之前不会返回。

    如果您不调用close(),根据实现/语言,您可能最终会出现资源/内存泄漏。

    【讨论】:

    • 我尝试过这样做,但现在的性能很糟糕,在 kafka 的最后一个版本中,有一个方法 send (List) 有什么办法可以用这个版本来做吗?跨度>
    • 不确定我是否可以跟随...您说“当我收到最后一条记录时,我关闭了 kafka 生产者”——因此,应该只有一次调用 .flush() 和一次调用to .close() -- 这会如何影响你的写作表现?
    • 我正在处理一个 xml,我正在使用 sax,当我检测到一个特殊标签时,我正在调用 flush 和 close 方法。但是如果我不这样做,性能会好很多,但是在一段时间后我会收到内存不足错误。在以前的 kafka 版本中,有一个发送消息列表的方法,但现在它不可用了。所以我正在考虑如何以更好的方式做到这一点
    • Producer内部缓存记录,分批发送。因此,不需要“发送列表”,因为它是由生产者自动处理的。你能描述一下你应用的整体模式吗?你什么时候创建生产者?你什么时候寄?你什么时候打电话关闭?我仍然不确定我是否完全理解您在做什么...还请查看文档:docs.confluent.io/current/clients/producer.html
    • 您可以使用单个生产者来处理整个 XML -- 无需为每个标签创建一个生产者。
    猜你喜欢
    • 1970-01-01
    • 2016-09-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-08-17
    • 2019-01-26
    • 1970-01-01
    相关资源
    最近更新 更多