Skip to content

Commit be19d13

Browse files
NaireenNaireen
and
Naireen
authored
add option to disable metrics (#34303)
Co-authored-by: Naireen <[email protected]>
1 parent 3828148 commit be19d13

File tree

1 file changed

+4
-1
lines changed

1 file changed

+4
-1
lines changed

sdks/java/io/kafka/src/main/java/org/apache/beam/sdk/io/kafka/KafkaIOInitializer.java

+4-1
Original file line numberDiff line numberDiff line change
@@ -19,13 +19,16 @@
1919

2020
import com.google.auto.service.AutoService;
2121
import org.apache.beam.sdk.harness.JvmInitializer;
22+
import org.apache.beam.sdk.options.ExperimentalOptions;
2223
import org.apache.beam.sdk.options.PipelineOptions;
2324

2425
/** Initialize KafkaIO feature flags on worker. */
2526
@AutoService(JvmInitializer.class)
2627
public class KafkaIOInitializer implements JvmInitializer {
2728
@Override
2829
public void beforeProcessing(PipelineOptions options) {
29-
KafkaSinkMetrics.setSupportKafkaMetrics(true);
30+
if (!ExperimentalOptions.hasExperiment(options, "disable_kafka_metrics")) {
31+
KafkaSinkMetrics.setSupportKafkaMetrics(true);
32+
}
3033
}
3134
}

0 commit comments

Comments
 (0)