【问题标题】:How can I create a cold publisher using Reactor-Kafka?如何使用 Reactor-Kafka 创建冷发布者?
【发布时间】:2017-12-02 13:54:31
【问题描述】:

是否可以在不订阅的情况下创建 Kafka 发布者,然后从另一个应用程序创建消费者,订阅主题并触发记录的发送?

我正在通过以下方式创建发布者:

  1. 致电KafkaSender.create(senderOptions)
  2. 后跟createOutbound()
  3. 只要应用程序正在运行,就会连续调用send()

在消费者方面(不同的应用程序),我所做的是:

  1. 致电KafkaReceiver.create(options)
  2. 后跟receive()
  3. 后跟subscribe(function -> doSomething())

目前,除非我在发布者上执行then().subscribe(),否则消费者不会收到任何信息,这会使其立即发出。理想情况下,我希望它在其他应用程序的消费者订阅时开始发射。

您能否告诉我我正在尝试做的事情是否可行?

非常感谢。

Reactor-Kafka 项目可以在这里找到:https://github.com/reactor/reactor-kafka

【问题讨论】:

  • 很高兴看到一些小应用程序让我们从我们身边玩。
  • 非常感谢您回复我。我创建了一些示例应用程序并在 BitBucket 上与您分享了代码。请注意,您需要在监控项目中编辑 reactive-consumer.propsreactive-producer.props 以指向真正的 Kafka 服务器。
  • 就目前而言,消费者接收来自热生产者的记录,而不是来自冷生产者的记录。如果您注释掉生产者的订阅,消费者将一无所获。提前感谢您帮助我解决这个问题!

标签: apache-kafka reactive-programming observer-pattern project-reactor


【解决方案1】:

嗯,我认为你应该从broker这个概念开始。这正是区分消费者和生产者的想法,让最后一个生产并忘记,真的不用担心另一边是否有消费者。

当消费者到达代理进行订阅时,它能够读取主题中的所有旧记录或仅对新生成的记录作出反应。但这已经是目标消息传递代理实现的细节。

绝对没有 API 可以让消费者找出话题上有消费者。恐怕您需要实施自己的解决方案。甚至可能通过同一个 Kafka 但使用单独的主题 - 例如发送/接收命令消息,并且已经存在,通过适当的命令,从生产者方面了解有一个生产者,因此调用提到的 then().subscribe()

【讨论】:

    猜你喜欢
    • 2023-01-31
    • 1970-01-01
    • 1970-01-01
    • 2021-04-11
    • 2019-08-25
    • 1970-01-01
    • 2021-12-13
    • 2017-10-02
    • 2020-08-03
    相关资源
    最近更新 更多