my custom flink metrics reporter don't work

Viewed 532

use flink version 1.13.1

i wrote a custom metrics reporter, but it's seems didn't work in my flink. when i start flink, the JobManager was show warn log like this:

2021-08-25 14:54:06,243 WARN  org.apache.flink.runtime.metrics.ReporterSetup               [] - The reporter factory (org.apache.flink.metrics.kafka.KafkaReporterFactory) could not be found for reporter kafka. Available factories: [org.apache.flink.metrics.slf4j.Slf4jReporterFactory, org.apache.flink.metrics.datadog.DatadogHttpReporterFactory, org.apache.flink.metrics.prometheus.PrometheusPushGatewayReporterFactory, org.apache.flink.metrics.graphite.GraphiteReporterFactory, org.apache.flink.metrics.statsd.StatsDReporterFactory, org.apache.flink.metrics.prometheus.PrometheusReporterFactory, org.apache.flink.metrics.jmx.JMXReporterFactory, org.apache.flink.metrics.influxdb.InfluxdbReporterFactory].
2021-08-25 14:54:06,245 INFO  org.apache.flink.runtime.metrics.MetricRegistryImpl          [] - No metrics reporter configured, no metrics will be exposed/reported.

but i already make folder which named 'metrics-kafka' in plugins folder than package the metrics reporter and copy the jar file to the 'metrics-kafka' or lib folder(both folder are didn't work)

my flink conf file:

metrics.reporter.kafka.factory.class: org.apache.flink.metrics.kafka.KafkaReporterFactory
metrics.reporter.kafka.class: org.apache.flink.metrics.kafka.KafkaReporter
metrics.reporter.kafka.interval: 15 SECONDS

my metrics reporter factory class:

package org.apache.flink.metrics.kafka

import org.apache.flink.metrics.reporter.{InterceptInstantiationViaReflection, MetricReporter, MetricReporterFactory}
import java.util.Properties

@InterceptInstantiationViaReflection(reporterClassName = "org.apache.flink.metrics.kafka.KafkaReporter")
class KafkaReporterFactory extends MetricReporterFactory{
  override def createMetricReporter(properties: Properties): MetricReporter = {
    new KafkaReporter()
  }
}

and my reporter class:

package org.apache.flink.metrics.kafka

import org.apache.flink.metrics.MetricConfig
import org.apache.flink.metrics.reporter.{InstantiateViaFactory, Scheduled}

@InstantiateViaFactory(factoryClassName = "org.apache.flink.metrics.kafka.KafkaReporterFactory")
class KafkaReporter extends MyAbstractReporter with Scheduled{
  some code ...
}
1 Answers

i fonud that need to add file in '/resources/META-INF/services/' which named 'org.apache.flink.metrics.reporter.MetricReporterFactory' and in this file write factory class path like 'org.apache.flink.metrics.kafka.KafkaReporterFactory'

Related