【问题标题】:How to use multiple counters in Flink如何在 Flink 中使用多个计数器
【发布时间】:2019-10-18 19:07:26
【问题描述】:

(有点像How to create dynamic metric in Flink

我有一个 events(someid:String, name:String) 流,出于监控原因,我需要一个计数器 per 事件 ID。 在所有 Flink 文档和示例中,我可以看到,例如,计数器是用 map 函数的 open 中的名称初始化的。

但在我的情况下,我无法初始化计数器,因为我需要每个 eventId 一个,而且我事先不知道该值。另外,我了解每次在 MapFunction 的 map() 方法中通过偶数时创建一个新计数器的成本是多么高。 最后,我不能保留计数器的“缓存”,因为它太大了。

理想情况下,我想要这样的东西:

class Event(id: String, name: String)

class ExampleMapFunction extends RichMapFunction[Event, Event] {
  @transient private var counter: Counter = _

  override def open(parameters: Configuration): Unit = {
    counter = new Counter()
  }

  override def map(event: Event): Event = {
    counter.inc(event.id)
    event
  }
}

或者基本上我可以实现我自己的计数器来允许我传递一个维度吗?如果是,怎么做?

对于这种用例有什么建议或最佳实践吗?

【问题讨论】:

  • 您能否解释一下为什么要为此使用指标而不是键控状态(这似乎是显而易见的答案)?指标并不能很好地扩展。
  • 出于监控原因,我想检查拓扑的每个步骤。例如,由于我有很多流连接,我想知道它不会加入的位置。

标签: scala apache-flink metrics


【解决方案1】:

如果保留计数器的缓存太大,那么我认为使用指标不会以满足您要求的方式进行扩展。

一些替代方案:

  • 使用辅助输出在一些外部、可查询/可视化的数据存储中收集有意义的事件 - 例如 influxdb。

  • 将信息保持在键控状态,并使用广播消息根据需要触发相关部分的输出(再次使用侧输出)。

  • 将信息保持在键控状态,并定期获取保存点,然后您可以使用状态处理器 API 通过查询对其进行分析。

【讨论】:

  • 感谢您的回答。当我说“它会太大”时,主要是因为它永远不会停止增加(自动生成的 id)。我将深入研究这些命题,但使用诸如 github.com/Netflix/spectator 之类的外部库怎么样?
猜你喜欢
  • 1970-01-01
  • 2020-11-22
  • 2022-12-15
  • 1970-01-01
  • 1970-01-01
  • 2018-02-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多