【发布时间】: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