Работа My Flink (1.6) прослушивает поток и выполняет некоторую агрегацию. Я хочу собрать показатели после агрегации, но у меня возникли некоторые трудности.
Мои метрики выглядят так:
id_1, 0.1
id_2, 0.3
...
Идентификаторы будут переменными, а значения будут увеличиваться и уменьшаться с течением времени, так что это выглядело как Датчик , который был наиболее подходящим.
Я создал эту функцию карты, чтобы фиксировать эти метрики в шкале:
class MetricsMapper extends RichMapFunction[MyObject, Double] {
override def map(obj: MyObject): Double = {
val metricVal = obj.metricVal
getRuntimeContext.getMetricGroup.gauge[Double, ScalaGauge[Double]](obj.id, ScalaGauge[Double](() => metricVal))
metricVal
}
}
Как это показывает, я использую свойство id моего объекта для регистрации датчика.
Проблема, с которой я столкнулся, заключается в том, что я получаю это предупреждение при запуске задания:
Name collision: Group already contains a Metric with the name "x" Metric will not be reported
Я интерпретирую это, поскольку мы уже создали этот датчик ранее в потоке, а новое значение игнорируется. Есть ли способ преодолеть это?
Спасибо