Skip to content

Commit 89f3db1

Browse files
committed
OTEL in spanner.
fdfd
1 parent aa8869c commit 89f3db1

16 files changed

Lines changed: 155 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
@@ -64,6 +64,7 @@ dependencies {
6464
implementation library.java.google_api_client
6565
implementation library.java.google_api_common
6666
implementation library.java.google_api_services_bigquery
67+
implementation library.java.opentelemetry_api
6768
implementation library.java.google_api_services_healthcare
6869
implementation library.java.google_api_services_pubsub
6970
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
@@ -39,6 +39,8 @@
3939
import com.google.spanner.v1.DirectedReadOptions;
4040
import com.google.spanner.v1.ExecuteSqlRequest;
4141
import com.google.spanner.v1.PartialResultSet;
42+
import io.opentelemetry.api.GlobalOpenTelemetry;
43+
import io.opentelemetry.api.OpenTelemetry;
4244
import java.util.HashSet;
4345
import java.util.Optional;
4446
import java.util.Set;
@@ -98,13 +100,17 @@ private SpannerAccessor(
98100
}
99101

100102
public static SpannerAccessor getOrCreate(SpannerConfig spannerConfig) {
103+
return getOrCreate(spannerConfig, GlobalOpenTelemetry.get());
104+
}
105+
106+
public static SpannerAccessor getOrCreate(SpannerConfig spannerConfig, OpenTelemetry otel) {
101107

102108
synchronized (spannerAccessors) {
103109
SpannerAccessor self = spannerAccessors.get(spannerConfig);
104110
if (self == null) {
105111
// Connect to spanner for this SpannerConfig.
106112
LOG.info("Connecting to {}", spannerConfig);
107-
self = SpannerAccessor.createAndConnect(spannerConfig);
113+
self = SpannerAccessor.createAndConnect(spannerConfig, otel);
108114
LOG.info("Successfully connected to {}", spannerConfig);
109115
spannerAccessors.put(spannerConfig, self);
110116
}
@@ -117,7 +123,27 @@ public static SpannerAccessor getOrCreate(SpannerConfig spannerConfig) {
117123

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

122148
// TODO(https://github.com/apache/beam/issues/37451) Disable gRPC gcp extension which was
123149
// causing the application thread to stall.
@@ -303,8 +329,8 @@ static SpannerOptions buildSpannerOptions(SpannerConfig spannerConfig) {
303329
return builder.build();
304330
}
305331

306-
private static SpannerAccessor createAndConnect(SpannerConfig spannerConfig) {
307-
SpannerOptions options = buildSpannerOptions(spannerConfig);
332+
private static SpannerAccessor createAndConnect(SpannerConfig spannerConfig, OpenTelemetry otel) {
333+
SpannerOptions options = buildSpannerOptions(spannerConfig, otel);
308334
Spanner spanner = options.getService();
309335
String instanceId = spannerConfig.getInstanceId().get();
310336
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
@@ -115,6 +115,8 @@ public String getHostValue() {
115115

116116
public abstract @Nullable ValueProvider<java.time.Duration> getWaitForSessionCreationDuration();
117117

118+
public abstract @Nullable ValueProvider<Boolean> getEnableOpenTelemetryTracing();
119+
118120
abstract Builder toBuilder();
119121

120122
public static SpannerConfig create() {
@@ -205,6 +207,9 @@ abstract Builder setExecuteStreamingSqlRetrySettings(
205207
abstract Builder setWaitForSessionCreationDuration(
206208
ValueProvider<java.time.Duration> waitForSessionCreationDuration);
207209

210+
abstract Builder setEnableOpenTelemetryTracing(
211+
ValueProvider<Boolean> enableOpenTelemetryTracing);
212+
208213
abstract Builder setClientCertPath(ValueProvider<String> clientCertPath);
209214

210215
abstract Builder setClientCertKeyPath(ValueProvider<String> clientCertKeyPath);
@@ -464,6 +469,16 @@ public SpannerConfig withWaitForSessionCreationDuration(
464469
ValueProvider.StaticValueProvider.of(waitForSessionCreationDuration));
465470
}
466471

472+
public SpannerConfig withEnableOpenTelemetryTracing(
473+
ValueProvider<Boolean> enableOpenTelemetryTracing) {
474+
return toBuilder().setEnableOpenTelemetryTracing(enableOpenTelemetryTracing).build();
475+
}
476+
477+
public SpannerConfig withEnableOpenTelemetryTracing(boolean enableOpenTelemetryTracing) {
478+
return withEnableOpenTelemetryTracing(
479+
ValueProvider.StaticValueProvider.of(enableOpenTelemetryTracing));
480+
}
481+
467482
/**
468483
* Specifies certificate paths to use for mTLS channel.
469484
*

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
@@ -62,6 +62,7 @@
6262
import com.google.gson.Gson;
6363
import com.google.gson.GsonBuilder;
6464
import com.google.spanner.v1.DirectedReadOptions;
65+
import io.opentelemetry.api.OpenTelemetry;
6566
import java.io.ByteArrayInputStream;
6667
import java.io.ByteArrayOutputStream;
6768
import java.io.IOException;
@@ -105,6 +106,7 @@
105106
import org.apache.beam.sdk.metrics.Lineage;
106107
import org.apache.beam.sdk.metrics.Metrics;
107108
import org.apache.beam.sdk.options.PipelineOptions;
109+
import org.apache.beam.sdk.options.SdkHarnessOptions;
108110
import org.apache.beam.sdk.options.StreamingOptions;
109111
import org.apache.beam.sdk.options.ValueProvider;
110112
import org.apache.beam.sdk.schemas.Schema;
@@ -1498,6 +1500,18 @@ public Write withHost(String host) {
14981500
return withHost(ValueProvider.StaticValueProvider.of(host));
14991501
}
15001502

1503+
/** Specifies whether OpenTelemetry tracing is enabled. */
1504+
public Write withEnableOpenTelemetryTracing(boolean enableOpenTelemetryTracing) {
1505+
return withEnableOpenTelemetryTracing(
1506+
ValueProvider.StaticValueProvider.of(enableOpenTelemetryTracing));
1507+
}
1508+
1509+
/** Specifies whether OpenTelemetry tracing is enabled. */
1510+
public Write withEnableOpenTelemetryTracing(ValueProvider<Boolean> enableOpenTelemetryTracing) {
1511+
SpannerConfig config = getSpannerConfig();
1512+
return withSpannerConfig(config.withEnableOpenTelemetryTracing(enableOpenTelemetryTracing));
1513+
}
1514+
15011515
/** Specifies the Cloud Spanner emulator host. */
15021516
public Write withEmulatorHost(ValueProvider<String> emulatorHost) {
15031517
SpannerConfig config = getSpannerConfig();
@@ -2024,6 +2038,19 @@ public ReadChangeStream withDatabaseId(ValueProvider<String> databaseId) {
20242038
return withSpannerConfig(config.withDatabaseId(databaseId));
20252039
}
20262040

2041+
/** Specifies whether OpenTelemetry tracing is enabled. */
2042+
public ReadChangeStream withEnableOpenTelemetryTracing(boolean enableOpenTelemetryTracing) {
2043+
return withEnableOpenTelemetryTracing(
2044+
ValueProvider.StaticValueProvider.of(enableOpenTelemetryTracing));
2045+
}
2046+
2047+
/** Specifies whether OpenTelemetry tracing is enabled. */
2048+
public ReadChangeStream withEnableOpenTelemetryTracing(
2049+
ValueProvider<Boolean> enableOpenTelemetryTracing) {
2050+
SpannerConfig config = getSpannerConfig();
2051+
return withSpannerConfig(config.withEnableOpenTelemetryTracing(enableOpenTelemetryTracing));
2052+
}
2053+
20272054
/** Specifies the change stream name. */
20282055
public ReadChangeStream withChangeStreamName(String changeStreamName) {
20292056
return toBuilder().setChangeStreamName(changeStreamName).build();
@@ -2272,7 +2299,9 @@ && getInclusiveStartAt().toSqlTimestamp().after(getInclusiveEndAt().toSqlTimesta
22722299
final ChangeStreamMetrics metrics = new ChangeStreamMetrics();
22732300
final RpcPriority rpcPriority = MoreObjects.firstNonNull(getRpcPriority(), RpcPriority.HIGH);
22742301
final SpannerAccessor spannerAccessor =
2275-
SpannerAccessor.getOrCreate(changeStreamSpannerConfig);
2302+
SpannerAccessor.getOrCreate(
2303+
changeStreamSpannerConfig,
2304+
input.getPipeline().getOptions().as(SdkHarnessOptions.class).getOpenTelemetry());
22762305
final boolean isMutableChangeStream =
22772306
isMutableChangeStream(
22782307
spannerAccessor.getDatabaseClient(), changeStreamDatabaseDialect, changeStreamName);
@@ -2413,7 +2442,11 @@ private static Dialect getDialect(SpannerConfig spannerConfig, PipelineOptions p
24132442
// Allow passing the credential from pipeline options to the getDialect() call.
24142443
SpannerConfig spannerConfigWithCredential =
24152444
buildSpannerConfigWithCredential(spannerConfig, pipelineOptions);
2416-
try (SpannerAccessor sa = SpannerAccessor.getOrCreate(spannerConfigWithCredential)) {
2445+
OpenTelemetry otel = null;
2446+
if (pipelineOptions != null) {
2447+
otel = pipelineOptions.as(SdkHarnessOptions.class).getOpenTelemetry();
2448+
}
2449+
try (SpannerAccessor sa = SpannerAccessor.getOrCreate(spannerConfigWithCredential, otel)) {
24172450
DatabaseClient databaseClient = sa.getDatabaseClient();
24182451
return databaseClient.getDialect();
24192452
}
@@ -2778,8 +2811,9 @@ static class WriteToSpannerFn extends DoFn<Iterable<MutationGroup>, Void> {
27782811
}
27792812

27802813
@Setup
2781-
public void setup() {
2782-
spannerAccessor = SpannerAccessor.getOrCreate(spannerConfig);
2814+
public void setup(PipelineOptions options) {
2815+
OpenTelemetry otel = options.as(SdkHarnessOptions.class).getOpenTelemetry();
2816+
spannerAccessor = SpannerAccessor.getOrCreate(spannerConfig, otel);
27832817
bundleWriteBackoff =
27842818
FluentBackoff.DEFAULT
27852819
.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: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,11 +20,13 @@
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;
2627
import org.apache.beam.sdk.io.gcp.spanner.SpannerConfig;
2728
import org.apache.beam.sdk.io.gcp.spanner.changestreams.ChangeStreamsConstants;
29+
import org.checkerframework.checker.nullness.qual.Nullable;
2830

2931
/**
3032
* Factory class to create data access objects to perform change stream queries and access the
@@ -41,6 +43,7 @@ public class DaoFactory implements Serializable {
4143
private transient PartitionMetadataAdminDao partitionMetadataAdminDao;
4244
private transient PartitionMetadataDao partitionMetadataDaoInstance;
4345
private transient ChangeStreamDao changeStreamDaoInstance;
46+
private transient @Nullable OpenTelemetry openTelemetry;
4447

4548
private final SpannerConfig changeStreamSpannerConfig;
4649
private final SpannerConfig metadataSpannerConfig;
@@ -100,6 +103,10 @@ public List<String> getTvfNameList() {
100103
return this.tvfNameList;
101104
}
102105

106+
public void setOpenTelemetry(@Nullable OpenTelemetry openTelemetry) {
107+
this.openTelemetry = openTelemetry;
108+
}
109+
103110
/**
104111
* Creates and returns a singleton DAO instance for admin operations over the partition metadata
105112
* table.
@@ -111,7 +118,8 @@ public List<String> getTvfNameList() {
111118
public synchronized PartitionMetadataAdminDao getPartitionMetadataAdminDao() {
112119
if (partitionMetadataAdminDao == null) {
113120
DatabaseAdminClient databaseAdminClient =
114-
SpannerAccessor.getOrCreate(metadataSpannerConfig).getDatabaseAdminClient();
121+
SpannerAccessor.getOrCreate(metadataSpannerConfig, this.openTelemetry)
122+
.getDatabaseAdminClient();
115123
partitionMetadataAdminDao =
116124
new PartitionMetadataAdminDao(
117125
databaseAdminClient,
@@ -131,7 +139,8 @@ public synchronized PartitionMetadataAdminDao getPartitionMetadataAdminDao() {
131139
* @return singleton instance of the {@link PartitionMetadataDao}
132140
*/
133141
public synchronized PartitionMetadataDao getPartitionMetadataDao() {
134-
final SpannerAccessor spannerAccessor = SpannerAccessor.getOrCreate(metadataSpannerConfig);
142+
final SpannerAccessor spannerAccessor =
143+
SpannerAccessor.getOrCreate(metadataSpannerConfig, this.openTelemetry);
135144
if (partitionMetadataDaoInstance == null) {
136145
partitionMetadataDaoInstance =
137146
new PartitionMetadataDao(
@@ -150,7 +159,8 @@ public synchronized PartitionMetadataDao getPartitionMetadataDao() {
150159
* @return singleton instance of the {@link ChangeStreamDao}
151160
*/
152161
public synchronized ChangeStreamDao getChangeStreamDao() {
153-
final SpannerAccessor spannerAccessor = SpannerAccessor.getOrCreate(changeStreamSpannerConfig);
162+
final SpannerAccessor spannerAccessor =
163+
SpannerAccessor.getOrCreate(changeStreamSpannerConfig, this.openTelemetry);
154164
if (changeStreamDaoInstance == null) {
155165
changeStreamDaoInstance =
156166
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)