【问题标题】:Flink - structuring job to maximize throughputFlink - 结构化作业以最大化吞吐量
【发布时间】:2015-12-01 22:59:41
【问题描述】:

我有 4 种类型的 kafka 主题,每种类型有 65 个主题。目标是对数据进行一些简单的窗口聚合并将其写入数据库。

拓扑如下所示:

kafka -> window -> reduce -> db write

在这个组合中的某个地方我想要/需要做一个联合 - 或者可能有几个(取决于每次合并多少主题)。

主题中的数据流范围从 10K 到 >200K 消息/分钟。

我有一个 30 核 / 节点的四节点 flink 集群。如何构建这些拓扑来分散负载?

【问题讨论】:

  • 快速提问,以确保并避免混淆:您总共有 260 个 Kafka 主题,每个主题都有自己的多个分区,还是每个 65 个分区有 4 个 Kafka 主题?在后一种情况下,散布自然会发生。
  • 260 个主题,每个主题有一个分区。

标签: java apache-flink flink-streaming


【解决方案1】:

我写这个答案是假设 65 个相同类型的主题中的每一个都包含相同类型的数据。

这个问题最简单的解决方案是更改 Kafka 设置,使您有 4 个主题,每个主题有 65 个分区。那么程序中有 4 个数据源,具有高并行度 (65),并且自然分布在整个集群中。

如果无法更改设置,我看到您可以做两件事:

  • 一种可能的解决方案是创建 FlinkKafkaConsumer 的修改版本,其中一个源可以使用多个主题(而不是一个主题的多个分区)。通过这种更改,它的工作方式几乎就像您使用许多分区而不是许多主题一样。如果你想使用这个解决方案,我会 ping 邮件列表以获得一些支持。无论如何,这将是对 Flink 代码的有价值的补充。

  • 您可以为每个源分配一个单独的资源组,这将为其分配一个专用插槽。你可以通过“env.addSource(new FlinkKafkaConsumer(...)).startNewResourceGroup();”来做到这一点。但是在这里,观察结果是您尝试在具有 120 个内核(因此可能有 120 个任务槽)的集群上执行 260 个不同的源。您需要增加槽数来容纳所有任务。

我认为第一个选项是更可取的选项。

【讨论】:

  • 所以如果我“按原样”运行,它会尝试将它们全部放在同一台机器上吗?无法更改分区。
  • 除非你指定“startNewResourceGroup()”,否则调度会尝试重用现有的资源组,这可能会导致它们共享同一台机器。在我看来,多主题 KafkaConsumer 是最好的选择。
  • 我将通过一个 FlinkKafkaConsumer issues.apache.org/jira/browse/FLINK-3102 添加对多个主题的读取支持
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2023-03-17
  • 1970-01-01
  • 2011-11-24
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-12-28
相关资源
最近更新 更多