Skip to content

Commit e012325

Browse files
committed
OTEL in spanner.
fdfd
1 parent 8d252c4 commit e012325

15 files changed

Lines changed: 153 additions & 30 deletions

File tree

sdks/java/io/google-cloud-platform/build.gradle

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,7 @@ dependencies {
5858
implementation library.java.google_api_client
5959
implementation library.java.google_api_common
6060
implementation library.java.google_api_services_bigquery
61+
implementation library.java.opentelemetry_api
6162
implementation library.java.google_api_services_healthcare
6263
implementation library.java.google_api_services_pubsub
6364
implementation library.java.google_api_services_storage

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/BatchSpannerRead.java

Lines changed: 9 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -27,13 +27,16 @@
2727
import com.google.cloud.spanner.SpannerException;
2828
import com.google.cloud.spanner.Struct;
2929
import com.google.cloud.spanner.TimestampBound;
30+
import io.opentelemetry.api.OpenTelemetry;
3031
import java.io.Serializable;
3132
import java.util.List;
3233
import java.util.Objects;
3334
import org.apache.beam.runners.core.metrics.ServiceCallMetric;
3435
import org.apache.beam.sdk.Pipeline;
3536
import org.apache.beam.sdk.io.gcp.spanner.SpannerIO.ReadAll;
3637
import org.apache.beam.sdk.metrics.Lineage;
38+
import org.apache.beam.sdk.options.PipelineOptions;
39+
import org.apache.beam.sdk.options.SdkHarnessOptions;
3740
import org.apache.beam.sdk.transforms.DoFn;
3841
import org.apache.beam.sdk.transforms.PTransform;
3942
import org.apache.beam.sdk.transforms.ParDo;
@@ -126,8 +129,9 @@ public GeneratePartitionsFn(
126129
}
127130

128131
@Setup
129-
public void setup() throws Exception {
130-
spannerAccessor = SpannerAccessor.getOrCreate(config);
132+
public void setup(PipelineOptions options) throws Exception {
133+
OpenTelemetry otel = options.as(SdkHarnessOptions.class).getOpenTelemetry();
134+
spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
131135
}
132136

133137
@Teardown
@@ -211,8 +215,9 @@ public ReadFromPartitionFn(
211215
}
212216

213217
@Setup
214-
public void setup() throws Exception {
215-
spannerAccessor = SpannerAccessor.getOrCreate(config);
218+
public void setup(PipelineOptions options) throws Exception {
219+
OpenTelemetry otel = options.as(SdkHarnessOptions.class).getOpenTelemetry();
220+
spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
216221

217222
// Use a LoadingCache for metrics as there can be different read operations which result in
218223
// different service call metrics labels. ServiceCallMetric items are created on-demand and

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/CreateTransactionFn.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,9 @@
1919

2020
import com.google.cloud.spanner.BatchReadOnlyTransaction;
2121
import com.google.cloud.spanner.TimestampBound;
22+
import io.opentelemetry.api.OpenTelemetry;
23+
import org.apache.beam.sdk.options.PipelineOptions;
24+
import org.apache.beam.sdk.options.SdkHarnessOptions;
2225
import org.apache.beam.sdk.transforms.DoFn;
2326

2427
/** Creates a batch transaction. */
@@ -37,8 +40,9 @@ class CreateTransactionFn extends DoFn<Object, Transaction> {
3740
private transient SpannerAccessor spannerAccessor;
3841

3942
@DoFn.Setup
40-
public void setup() throws Exception {
41-
spannerAccessor = SpannerAccessor.getOrCreate(config);
43+
public void setup(PipelineOptions options) throws Exception {
44+
OpenTelemetry otel = options.as(SdkHarnessOptions.class).getOpenTelemetry();
45+
spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
4246
}
4347

4448
@Teardown

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/NaiveSpannerRead.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,10 +25,13 @@
2525
import com.google.cloud.spanner.SpannerException;
2626
import com.google.cloud.spanner.Struct;
2727
import com.google.cloud.spanner.TimestampBound;
28+
import io.opentelemetry.api.OpenTelemetry;
2829
import java.util.Objects;
2930
import org.apache.beam.runners.core.metrics.ServiceCallMetric;
3031
import org.apache.beam.sdk.Pipeline;
3132
import org.apache.beam.sdk.metrics.Lineage;
33+
import org.apache.beam.sdk.options.PipelineOptions;
34+
import org.apache.beam.sdk.options.SdkHarnessOptions;
3235
import org.apache.beam.sdk.transforms.DoFn;
3336
import org.apache.beam.sdk.transforms.PTransform;
3437
import org.apache.beam.sdk.transforms.ParDo;
@@ -97,8 +100,9 @@ private static class NaiveSpannerReadFn extends DoFn<ReadOperation, Struct> {
97100
}
98101

99102
@Setup
100-
public void setup() throws Exception {
101-
spannerAccessor = SpannerAccessor.getOrCreate(config);
103+
public void setup(PipelineOptions options) throws Exception {
104+
OpenTelemetry otel = options.as(SdkHarnessOptions.class).getOpenTelemetry();
105+
spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
102106
projectId = SpannerIO.resolveSpannerProjectId(config);
103107
}
104108

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/ReadSpannerSchema.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,8 +22,11 @@
2222
import com.google.cloud.spanner.ReadOnlyTransaction;
2323
import com.google.cloud.spanner.ResultSet;
2424
import com.google.cloud.spanner.Statement;
25+
import io.opentelemetry.api.OpenTelemetry;
2526
import java.util.HashSet;
2627
import java.util.Set;
28+
import org.apache.beam.sdk.options.PipelineOptions;
29+
import org.apache.beam.sdk.options.SdkHarnessOptions;
2730
import org.apache.beam.sdk.transforms.DoFn;
2831
import org.apache.beam.sdk.values.PCollectionView;
2932

@@ -75,8 +78,9 @@ public ReadSpannerSchema(
7578
}
7679

7780
@Setup
78-
public void setup() throws Exception {
79-
spannerAccessor = SpannerAccessor.getOrCreate(config);
81+
public void setup(PipelineOptions options) throws Exception {
82+
OpenTelemetry otel = options.as(SdkHarnessOptions.class).getOpenTelemetry();
83+
spannerAccessor = SpannerAccessor.getOrCreate(config, otel);
8084
}
8185

8286
@Teardown

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerAccessor.java

Lines changed: 29 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,8 @@
3838
import com.google.spanner.v1.CommitResponse;
3939
import com.google.spanner.v1.ExecuteSqlRequest;
4040
import com.google.spanner.v1.PartialResultSet;
41+
import io.opentelemetry.api.GlobalOpenTelemetry;
42+
import io.opentelemetry.api.OpenTelemetry;
4143
import java.util.HashSet;
4244
import java.util.Optional;
4345
import java.util.Set;
@@ -97,13 +99,17 @@ private SpannerAccessor(
9799
}
98100

99101
public static SpannerAccessor getOrCreate(SpannerConfig spannerConfig) {
102+
return getOrCreate(spannerConfig, GlobalOpenTelemetry.get());
103+
}
104+
105+
public static SpannerAccessor getOrCreate(SpannerConfig spannerConfig, OpenTelemetry otel) {
100106

101107
synchronized (spannerAccessors) {
102108
SpannerAccessor self = spannerAccessors.get(spannerConfig);
103109
if (self == null) {
104110
// Connect to spanner for this SpannerConfig.
105111
LOG.info("Connecting to {}", spannerConfig);
106-
self = SpannerAccessor.createAndConnect(spannerConfig);
112+
self = SpannerAccessor.createAndConnect(spannerConfig, otel);
107113
LOG.info("Successfully connected to {}", spannerConfig);
108114
spannerAccessors.put(spannerConfig, self);
109115
}
@@ -116,7 +122,27 @@ public static SpannerAccessor getOrCreate(SpannerConfig spannerConfig) {
116122

117123
@VisibleForTesting
118124
static SpannerOptions buildSpannerOptions(SpannerConfig spannerConfig) {
125+
return buildSpannerOptions(spannerConfig, GlobalOpenTelemetry.get());
126+
}
127+
128+
@VisibleForTesting
129+
static SpannerOptions buildSpannerOptions(SpannerConfig spannerConfig, OpenTelemetry otel) {
119130
SpannerOptions.Builder builder = SpannerOptions.newBuilder();
131+
if (otel != null) {
132+
builder.setOpenTelemetry(otel);
133+
}
134+
ValueProvider<Boolean> enableOpenTelemetryTracing =
135+
spannerConfig.getEnableOpenTelemetryTracing();
136+
if (enableOpenTelemetryTracing != null
137+
&& enableOpenTelemetryTracing.isAccessible()
138+
&& enableOpenTelemetryTracing.get()) {
139+
builder.setEnableExtendedTracing(true);
140+
builder.setEnableEndToEndTracing(true);
141+
builder.setEnableApiTracing(true);
142+
SpannerOptions.disableOpenCensusMetrics();
143+
SpannerOptions.enableOpenTelemetryMetrics();
144+
SpannerOptions.enableOpenTelemetryTraces();
145+
}
120146

121147
// TODO(https://github.com/apache/beam/issues/37451) Disable gRPC gcp extension which was
122148
// causing the application thread to stall.
@@ -298,8 +324,8 @@ static SpannerOptions buildSpannerOptions(SpannerConfig spannerConfig) {
298324
return builder.build();
299325
}
300326

301-
private static SpannerAccessor createAndConnect(SpannerConfig spannerConfig) {
302-
SpannerOptions options = buildSpannerOptions(spannerConfig);
327+
private static SpannerAccessor createAndConnect(SpannerConfig spannerConfig, OpenTelemetry otel) {
328+
SpannerOptions options = buildSpannerOptions(spannerConfig, otel);
303329
Spanner spanner = options.getService();
304330
String instanceId = spannerConfig.getInstanceId().get();
305331
String databaseId = spannerConfig.getDatabaseId().get();

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerConfig.java

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -110,6 +110,8 @@ public String getHostValue() {
110110

111111
public abstract @Nullable ValueProvider<java.time.Duration> getWaitForSessionCreationDuration();
112112

113+
public abstract @Nullable ValueProvider<Boolean> getEnableOpenTelemetryTracing();
114+
113115
abstract Builder toBuilder();
114116

115117
public static SpannerConfig create() {
@@ -198,6 +200,9 @@ abstract Builder setExecuteStreamingSqlRetrySettings(
198200
abstract Builder setWaitForSessionCreationDuration(
199201
ValueProvider<java.time.Duration> waitForSessionCreationDuration);
200202

203+
abstract Builder setEnableOpenTelemetryTracing(
204+
ValueProvider<Boolean> enableOpenTelemetryTracing);
205+
201206
abstract Builder setClientCertPath(ValueProvider<String> clientCertPath);
202207

203208
abstract Builder setClientCertKeyPath(ValueProvider<String> clientCertKeyPath);
@@ -423,6 +428,16 @@ public SpannerConfig withWaitForSessionCreationDuration(
423428
ValueProvider.StaticValueProvider.of(waitForSessionCreationDuration));
424429
}
425430

431+
public SpannerConfig withEnableOpenTelemetryTracing(
432+
ValueProvider<Boolean> enableOpenTelemetryTracing) {
433+
return toBuilder().setEnableOpenTelemetryTracing(enableOpenTelemetryTracing).build();
434+
}
435+
436+
public SpannerConfig withEnableOpenTelemetryTracing(boolean enableOpenTelemetryTracing) {
437+
return withEnableOpenTelemetryTracing(
438+
ValueProvider.StaticValueProvider.of(enableOpenTelemetryTracing));
439+
}
440+
426441
/**
427442
* Specifies certificate paths to use for mTLS channel.
428443
*

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/SpannerIO.java

Lines changed: 38 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,7 @@
6161
import com.google.cloud.spanner.TimestampBound;
6262
import com.google.gson.Gson;
6363
import com.google.gson.GsonBuilder;
64+
import io.opentelemetry.api.OpenTelemetry;
6465
import java.io.ByteArrayInputStream;
6566
import java.io.ByteArrayOutputStream;
6667
import java.io.IOException;
@@ -104,6 +105,7 @@
104105
import org.apache.beam.sdk.metrics.Lineage;
105106
import org.apache.beam.sdk.metrics.Metrics;
106107
import org.apache.beam.sdk.options.PipelineOptions;
108+
import org.apache.beam.sdk.options.SdkHarnessOptions;
107109
import org.apache.beam.sdk.options.StreamingOptions;
108110
import org.apache.beam.sdk.options.ValueProvider;
109111
import org.apache.beam.sdk.schemas.Schema;
@@ -1461,6 +1463,18 @@ public Write withHost(String host) {
14611463
return withHost(ValueProvider.StaticValueProvider.of(host));
14621464
}
14631465

1466+
/** Specifies whether OpenTelemetry tracing is enabled. */
1467+
public Write withEnableOpenTelemetryTracing(boolean enableOpenTelemetryTracing) {
1468+
return withEnableOpenTelemetryTracing(
1469+
ValueProvider.StaticValueProvider.of(enableOpenTelemetryTracing));
1470+
}
1471+
1472+
/** Specifies whether OpenTelemetry tracing is enabled. */
1473+
public Write withEnableOpenTelemetryTracing(ValueProvider<Boolean> enableOpenTelemetryTracing) {
1474+
SpannerConfig config = getSpannerConfig();
1475+
return withSpannerConfig(config.withEnableOpenTelemetryTracing(enableOpenTelemetryTracing));
1476+
}
1477+
14641478
/** Specifies the Cloud Spanner emulator host. */
14651479
public Write withEmulatorHost(ValueProvider<String> emulatorHost) {
14661480
SpannerConfig config = getSpannerConfig();
@@ -1987,6 +2001,19 @@ public ReadChangeStream withDatabaseId(ValueProvider<String> databaseId) {
19872001
return withSpannerConfig(config.withDatabaseId(databaseId));
19882002
}
19892003

2004+
/** Specifies whether OpenTelemetry tracing is enabled. */
2005+
public ReadChangeStream withEnableOpenTelemetryTracing(boolean enableOpenTelemetryTracing) {
2006+
return withEnableOpenTelemetryTracing(
2007+
ValueProvider.StaticValueProvider.of(enableOpenTelemetryTracing));
2008+
}
2009+
2010+
/** Specifies whether OpenTelemetry tracing is enabled. */
2011+
public ReadChangeStream withEnableOpenTelemetryTracing(
2012+
ValueProvider<Boolean> enableOpenTelemetryTracing) {
2013+
SpannerConfig config = getSpannerConfig();
2014+
return withSpannerConfig(config.withEnableOpenTelemetryTracing(enableOpenTelemetryTracing));
2015+
}
2016+
19902017
/** Specifies the change stream name. */
19912018
public ReadChangeStream withChangeStreamName(String changeStreamName) {
19922019
return toBuilder().setChangeStreamName(changeStreamName).build();
@@ -2213,7 +2240,9 @@ && getInclusiveStartAt().toSqlTimestamp().after(getInclusiveEndAt().toSqlTimesta
22132240
final ChangeStreamMetrics metrics = new ChangeStreamMetrics();
22142241
final RpcPriority rpcPriority = MoreObjects.firstNonNull(getRpcPriority(), RpcPriority.HIGH);
22152242
final SpannerAccessor spannerAccessor =
2216-
SpannerAccessor.getOrCreate(changeStreamSpannerConfig);
2243+
SpannerAccessor.getOrCreate(
2244+
changeStreamSpannerConfig,
2245+
input.getPipeline().getOptions().as(SdkHarnessOptions.class).getOpenTelemetry());
22172246
final boolean isMutableChangeStream =
22182247
isMutableChangeStream(
22192248
spannerAccessor.getDatabaseClient(), changeStreamDatabaseDialect, changeStreamName);
@@ -2354,7 +2383,11 @@ private static Dialect getDialect(SpannerConfig spannerConfig, PipelineOptions p
23542383
// Allow passing the credential from pipeline options to the getDialect() call.
23552384
SpannerConfig spannerConfigWithCredential =
23562385
buildSpannerConfigWithCredential(spannerConfig, pipelineOptions);
2357-
try (SpannerAccessor sa = SpannerAccessor.getOrCreate(spannerConfigWithCredential)) {
2386+
OpenTelemetry otel = null;
2387+
if (pipelineOptions != null) {
2388+
otel = pipelineOptions.as(SdkHarnessOptions.class).getOpenTelemetry();
2389+
}
2390+
try (SpannerAccessor sa = SpannerAccessor.getOrCreate(spannerConfigWithCredential, otel)) {
23582391
DatabaseClient databaseClient = sa.getDatabaseClient();
23592392
return databaseClient.getDialect();
23602393
}
@@ -2719,8 +2752,9 @@ static class WriteToSpannerFn extends DoFn<Iterable<MutationGroup>, Void> {
27192752
}
27202753

27212754
@Setup
2722-
public void setup() {
2723-
spannerAccessor = SpannerAccessor.getOrCreate(spannerConfig);
2755+
public void setup(PipelineOptions options) {
2756+
OpenTelemetry otel = options.as(SdkHarnessOptions.class).getOpenTelemetry();
2757+
spannerAccessor = SpannerAccessor.getOrCreate(spannerConfig, otel);
27242758
bundleWriteBackoff =
27252759
FluentBackoff.DEFAULT
27262760
.withMaxCumulativeBackoff(spannerConfig.getMaxCumulativeBackoff().get())

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/DaoFactory.java

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
import com.google.cloud.spanner.DatabaseAdminClient;
2121
import com.google.cloud.spanner.Dialect;
2222
import com.google.cloud.spanner.Options.RpcPriority;
23+
import io.opentelemetry.api.OpenTelemetry;
2324
import java.io.Serializable;
2425
import java.util.List;
2526
import org.apache.beam.sdk.io.gcp.spanner.SpannerAccessor;
@@ -41,6 +42,7 @@ public class DaoFactory implements Serializable {
4142
private transient PartitionMetadataAdminDao partitionMetadataAdminDao;
4243
private transient PartitionMetadataDao partitionMetadataDaoInstance;
4344
private transient ChangeStreamDao changeStreamDaoInstance;
45+
private transient OpenTelemetry openTelemetry;
4446

4547
private final SpannerConfig changeStreamSpannerConfig;
4648
private final SpannerConfig metadataSpannerConfig;
@@ -100,6 +102,10 @@ public List<String> getTvfNameList() {
100102
return this.tvfNameList;
101103
}
102104

105+
public void setOpenTelemetry(OpenTelemetry openTelemetry) {
106+
this.openTelemetry = openTelemetry;
107+
}
108+
103109
/**
104110
* Creates and returns a singleton DAO instance for admin operations over the partition metadata
105111
* table.
@@ -111,7 +117,8 @@ public List<String> getTvfNameList() {
111117
public synchronized PartitionMetadataAdminDao getPartitionMetadataAdminDao() {
112118
if (partitionMetadataAdminDao == null) {
113119
DatabaseAdminClient databaseAdminClient =
114-
SpannerAccessor.getOrCreate(metadataSpannerConfig).getDatabaseAdminClient();
120+
SpannerAccessor.getOrCreate(metadataSpannerConfig, this.openTelemetry)
121+
.getDatabaseAdminClient();
115122
partitionMetadataAdminDao =
116123
new PartitionMetadataAdminDao(
117124
databaseAdminClient,
@@ -131,7 +138,8 @@ public synchronized PartitionMetadataAdminDao getPartitionMetadataAdminDao() {
131138
* @return singleton instance of the {@link PartitionMetadataDao}
132139
*/
133140
public synchronized PartitionMetadataDao getPartitionMetadataDao() {
134-
final SpannerAccessor spannerAccessor = SpannerAccessor.getOrCreate(metadataSpannerConfig);
141+
final SpannerAccessor spannerAccessor =
142+
SpannerAccessor.getOrCreate(metadataSpannerConfig, this.openTelemetry);
135143
if (partitionMetadataDaoInstance == null) {
136144
partitionMetadataDaoInstance =
137145
new PartitionMetadataDao(
@@ -150,7 +158,8 @@ public synchronized PartitionMetadataDao getPartitionMetadataDao() {
150158
* @return singleton instance of the {@link ChangeStreamDao}
151159
*/
152160
public synchronized ChangeStreamDao getChangeStreamDao() {
153-
final SpannerAccessor spannerAccessor = SpannerAccessor.getOrCreate(changeStreamSpannerConfig);
161+
final SpannerAccessor spannerAccessor =
162+
SpannerAccessor.getOrCreate(changeStreamSpannerConfig, this.openTelemetry);
154163
if (changeStreamDaoInstance == null) {
155164
changeStreamDaoInstance =
156165
new ChangeStreamDao(

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/CleanUpReadChangeStreamDoFn.java

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,8 @@
2020
import java.io.Serializable;
2121
import java.util.List;
2222
import org.apache.beam.sdk.io.gcp.spanner.changestreams.dao.DaoFactory;
23+
import org.apache.beam.sdk.options.PipelineOptions;
24+
import org.apache.beam.sdk.options.SdkHarnessOptions;
2325
import org.apache.beam.sdk.transforms.DoFn;
2426

2527
public class CleanUpReadChangeStreamDoFn extends DoFn<byte[], Void> implements Serializable {
@@ -32,6 +34,11 @@ public CleanUpReadChangeStreamDoFn(DaoFactory daoFactory) {
3234
this.daoFactory = daoFactory;
3335
}
3436

37+
@Setup
38+
public void setup(PipelineOptions options) {
39+
daoFactory.setOpenTelemetry(options.as(SdkHarnessOptions.class).getOpenTelemetry());
40+
}
41+
3542
@ProcessElement
3643
public void processElement(OutputReceiver<Void> receiver) {
3744
List<String> indexes = daoFactory.getPartitionMetadataDao().findAllTableIndexes();

0 commit comments

Comments
 (0)