【问题标题】:How to set TOPOLOGY_MAX_SPOUT_PENDING parameter如何设置 TOPOLOGY_MAX_SPOUT_PENDING 参数
【发布时间】:2015-12-25 01:44:53
【问题描述】:

在我的拓扑中,我从 Kafka 队列中读取触发消息。收到触发消息后,我需要向螺栓发出大约 4096 条消息。在 Bolt 中,经过一些处理后,它将发布到另一个 Kafka 队列(稍后另一个拓扑会使用它)。

我正在尝试设置TOPOLOGY_MAX_SPOUT_PENDING 参数来限制要发送的消息数量。但我看到它没有任何效果。是因为我在一个nextTuple() 方法中发出所有元组吗?如果是这样,应该如何解决?

【问题讨论】:

  • 你试过什么代码?
  • 我已经编辑了您的帖子以包含格式标签,并且还修复了一些拼写错误。问题越清楚,答案就越好!

标签: java apache-kafka apache-storm


【解决方案1】:

如果你从 kafka 阅读,你应该使用 Storm 自带的KafkaSpout。不要尝试实现自己的 spout,相信我,我在生产中使用 KafkaSpout,它工作得非常顺利。每个 Kafka 消息只生成一个元组。

正如您在this nice page from the manual 上看到的,您可以像这样设置topology.max.spout.pending

Config conf = new Config();
conf.setMaxSpoutPending(5000);
StormSubmitter.submitTopology("mytopology", conf, topology);

topology.max.spout.pending 是为每个 spout 设置的,如果您有四个 spout,您的拓扑中的不完整元组的最大值将等于 spout 的数量 * topology.max.spout.pending。

另一个提示是,您应该使用storm UI 来查看topology.max.spout.pending 是否设置正确。


记住topology.max.spout.pending只是拓扑内部未处理的元组的数量,拓扑永远不会停止消费来自kafka的消息,至少在生产系统上......如果你想批量消费4096你需要在你的 bolts 上实现缓存逻辑,或者使用 Storm 以外的东西(面向微批处理的东西)。

【讨论】:

  • 感谢您的回复。
【解决方案2】:

要使 TOPOLOGY_MAX_SPOUT_PENDING 生效,您需要启用容错机制(即,在 Spout 中分配消息 ID,在 Bolts 中分配锚点和确认)。此外,如果每次调用 Spout.nextTuple() TOPOLOGY_MAX_SPOUT_PENDING 发出多个元组,则将无法按预期工作。

由于其他一些原因,这实际上是一种不好的做法,因此每次Spout.nextTuple() 调用发出多个元组(有关更多详细信息,请参阅Why should I not loop or block in Spout.nextTuple())。

【讨论】:

  • 感谢马蒂亚斯的回复。我在这里看到了问题,因为我为一条消息发出了所有 4096 元组。但是,这是我的用例要求我做的事情。
  • 也许您可以将 TOPOLOGY_MAX_SPOUT_PENDING 设置为 1。这应该会触发对 nextTuple() 的一次调用,并且在您发出的所有 4096 个元组都得到处理之前不会发出第二次调用。
  • 感谢马泰斯的回复。我会重新考虑我的拓扑设计。
猜你喜欢
  • 1970-01-01
  • 2013-11-30
  • 2019-04-10
  • 2012-09-17
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多