We read every piece of feedback, and take your input very seriously.
To see all available qualifiers, see our documentation.
1 parent 3828148 commit be19d13Copy full SHA for be19d13
sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIOInitializer.java
@@ -19,13 +19,16 @@
19
20
import com.google.auto.service.AutoService;
21
import org.apache.beam.sdk.harness.JvmInitializer;
22
+import org.apache.beam.sdk.options.ExperimentalOptions;
23
import org.apache.beam.sdk.options.PipelineOptions;
24
25
/** Initialize KafkaIO feature flags on worker. */
26
@AutoService(JvmInitializer.class)
27
public class KafkaIOInitializer implements JvmInitializer {
28
@Override
29
public void beforeProcessing(PipelineOptions options) {
- KafkaSinkMetrics.setSupportKafkaMetrics(true);
30
+ if (!ExperimentalOptions.hasExperiment(options, "disable_kafka_metrics")) {
31
+ KafkaSinkMetrics.setSupportKafkaMetrics(true);
32
+ }
33
}
34
0 commit comments