【问题标题】:Buffering in a Windowed Kafka Streams App在窗口化的 Kafka Streams 应用程序中进行缓冲
【发布时间】:2018-12-19 23:10:42
【问题描述】:

在我们的应用程序中,我们试图从输入主题中获取 JSON 消息,将它们组合到给定窗口中,然后将它们写入目标主题。 mergeJsonNodes 是负责简单合并两个 JSON 对象的函数。

KStream<String, JsonNode> transformed = datastreamSource
  .groupByKey(Serialized.with(Serdes.String(), JSON_SERDE))
  .windowedBy(SessionWindows.with(60 * 1000))
  .reduce((a, b) -> mergeJsonNodes(a, b))
  .toStream((windowedKey, node) -> windowedKey.key());

我们已经在我们的几个非生产环境中成功地部署了它。然而,当我们转向生产时,输入主题 (datastreamSource) 的数量要大得多,我们遇到了一个我们正在寻求理解的瓶颈。

我们看到的是,我们的流应用程序在源主题上进展缓慢,并且正在致力于目标主题〜每分钟一次。但是,它从输入主题中摄取的速度太慢,无法跟上我们致力于该主题的生产流量。我们正在从一个几个月来一直表现良好的非窗口化、非分组流应用程序迁移。

Kafka 流应用程序的主机上的资源似乎没有受到限制;不是应用程序缺少内存或磁盘。

我们的问题是我们可以修改哪些其他因素,特别是配置设置,以允许流应用一次从输入主题中提取更多消息。看来我们的应用程序在继续从源代码顶部阅读的能力方面受到了某种限制 我知道了。

最初从the docs 跳出来的两个:
* buffered.records.per.partition
* cache.max.bytes.buffering

有没有人使用过可以提供任何指示的高吞吐量窗口流应用程序?谢谢!!

【问题讨论】:

  • 你应该弄清楚瓶颈是什么。你的目标负载是多少?你的网络饱和了吗?你有多少个输入主题分区?你运行了多少线程/实例?

标签: apache-kafka streaming apache-kafka-streams


【解决方案1】:

我不知道特别是在窗口聚合中,但是在 Kafka 流中进行聚合时,您有 2 个配置来查看处理聚合处理器节点在刷新到状态存储并将结果聚合记录发送到下游处理器之前如何缓存消息: cache.max.bytes.buffering, commit.interval.ms.

您拥有可以在 kafka 流中调整的消费者配置:poll.ms

您也可以扩展您的应用程序,您的输入主题有多少个分区?它将导致处理您输入主题的任务数量,从而影响您应用的可扩展性。

更多的分区意味着更多的任务意味着更多的消费者意味着更多的实例,甚至更多的实例线程,(检查num.streams.thread)。

希望对您有所帮助。

【讨论】:

    猜你喜欢
    • 2018-11-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-01-23
    • 1970-01-01
    • 2018-08-21
    • 2018-07-17
    • 2016-12-27
    相关资源
    最近更新 更多