Skip to content

Commit fcdd87f

Browse files
authored
[OpenTelemetry] FNHarness, Dataflow Runner v1 - set open telemetry settings. (#38785)
1 parent a0c9606 commit fcdd87f

2 files changed

Lines changed: 24 additions & 0 deletions

File tree

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -115,6 +115,7 @@
115115
import org.apache.beam.sdk.io.gcp.bigquery.BigQuerySinkMetrics;
116116
import org.apache.beam.sdk.metrics.MetricsEnvironment;
117117
import org.apache.beam.sdk.options.ExperimentalOptions;
118+
import org.apache.beam.sdk.options.SdkHarnessOptions;
118119
import org.apache.beam.sdk.util.construction.CoderTranslation;
119120
import org.apache.beam.sdk.values.WindowedValues;
120121
import org.apache.beam.vendor.grpc.v1p69p0.io.grpc.auth.MoreCallCredentials;
@@ -1042,6 +1043,18 @@ public static void main(String[] args) throws Exception {
10421043
WindowedValues.FullWindowedValueCoder.setMetadataSupported();
10431044
}
10441045

1046+
SdkHarnessOptions sdkHarnessOptions = options.as(SdkHarnessOptions.class);
1047+
Map<String, String> openTelemetryProperties = sdkHarnessOptions.getOpenTelemetryProperties();
1048+
if (openTelemetryProperties != null && !openTelemetryProperties.isEmpty()) {
1049+
openTelemetryProperties.forEach(
1050+
(k, v) -> {
1051+
if (k != null && v != null) {
1052+
System.setProperty(k, v);
1053+
}
1054+
});
1055+
LOG.info("Enabled Open Telemetry with properties: {}", openTelemetryProperties);
1056+
}
1057+
10451058
LOG.debug("Creating StreamingDataflowWorker from options: {}", options);
10461059
StreamingDataflowWorker worker = StreamingDataflowWorker.fromOptions(options);
10471060

sdks/java/harness/src/main/java/org/apache/beam/fn/harness/FnHarness.java

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -296,6 +296,17 @@ public static void main(
296296
// Register standard file systems.
297297
FileSystems.setDefaultPipelineOptions(options);
298298
CoderTranslation.verifyModelCodersRegistered();
299+
SdkHarnessOptions sdkHarnessOptions = options.as(SdkHarnessOptions.class);
300+
Map<String, String> openTelemetryProperties = sdkHarnessOptions.getOpenTelemetryProperties();
301+
if (openTelemetryProperties != null && !openTelemetryProperties.isEmpty()) {
302+
openTelemetryProperties.forEach(
303+
(k, v) -> {
304+
if (k != null && v != null) {
305+
System.setProperty(k, v);
306+
}
307+
});
308+
LOG.info("Enabled Open Telemetry with properties: {}", openTelemetryProperties);
309+
}
299310
EnumMap<
300311
BeamFnApi.InstructionRequest.RequestCase,
301312
ThrowingFunction<InstructionRequest, BeamFnApi.InstructionResponse.Builder>>

0 commit comments

Comments
 (0)