Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions sdks/java/io/solace/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -27,9 +27,12 @@ description = "Apache Beam :: SDKs :: Java :: IO :: Solace"
ext.summary = """IO to read and write to Solace destinations (queues and topics)."""

dependencies {
implementation platform(library.java.opentelemetry_bom)
implementation enforcedPlatform(library.java.google_cloud_platform_libraries_bom)
implementation library.java.vendored_guava_32_1_2_jre
implementation project(path: ":sdks:java:core", configuration: "shadow")
implementation library.java.opentelemetry_api
implementation library.java.opentelemetry_context
implementation library.java.slf4j_api
implementation library.java.joda_time
implementation library.java.solace
Expand All @@ -51,6 +54,8 @@ dependencies {
implementation library.java.jackson_databind

testImplementation library.java.junit
testImplementation library.java.opentelemetry_sdk
testImplementation library.java.opentelemetry_sdk_testing
testImplementation project(path: ":sdks:java:io:common")
testImplementation project(path: ":sdks:java:testing:test-utils")
testImplementation project(path: ":sdks:java:core", configuration: "shadowTest")
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,18 @@
import com.solacesystems.jcsmp.JCSMPFactory;
import com.solacesystems.jcsmp.Queue;
import com.solacesystems.jcsmp.Topic;
import io.opentelemetry.api.OpenTelemetry;
import io.opentelemetry.api.trace.Span;
import io.opentelemetry.api.trace.SpanKind;
import io.opentelemetry.api.trace.Tracer;
import io.opentelemetry.api.trace.propagation.W3CTraceContextPropagator;
import io.opentelemetry.context.Context;
import io.opentelemetry.context.Scope;
import io.opentelemetry.context.propagation.TextMapGetter;
import io.opentelemetry.context.propagation.TextMapSetter;
import java.io.IOException;
import java.util.HashMap;
import java.util.Map;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.annotations.Internal;
import org.apache.beam.sdk.coders.CannotProvideCoderException;
Expand All @@ -47,7 +58,9 @@
import org.apache.beam.sdk.io.solace.write.UnboundedBatchedSolaceWriter;
import org.apache.beam.sdk.io.solace.write.UnboundedStreamingSolaceWriter;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.schemas.NoSuchSchemaException;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.transforms.PTransform;
import org.apache.beam.sdk.transforms.ParDo;
Expand Down Expand Up @@ -423,6 +436,7 @@ public class SolaceIO {
Duration.standardSeconds(30);
private static final Duration DEFAULT_ACK_DEADLINE = Duration.standardSeconds(30);
public static final boolean DEFAULT_NACK_ON_TIMEOUT = false;
public static final boolean DEFAULT_ENABLE_OPENTELEMETRY_TRACING = false;
public static final int DEFAULT_WRITER_NUM_SHARDS = 20;
public static final int DEFAULT_WRITER_CLIENTS_PER_WORKER = 4;
public static final Boolean DEFAULT_WRITER_PUBLISH_LATENCY_METRICS = false;
Expand Down Expand Up @@ -472,7 +486,8 @@ public static Read<Solace.Record> read() {
.setDeduplicateRecords(DEFAULT_DEDUPLICATE_RECORDS)
.setWatermarkIdleDurationThreshold(DEFAULT_WATERMARK_IDLE_DURATION_THRESHOLD)
.setAckDeadline(DEFAULT_ACK_DEADLINE)
.setNackOnTimeout(DEFAULT_NACK_ON_TIMEOUT));
.setNackOnTimeout(DEFAULT_NACK_ON_TIMEOUT)
.setEnableOpenTelemetryTracing(DEFAULT_ENABLE_OPENTELEMETRY_TRACING));
}

/**
Expand Down Expand Up @@ -503,7 +518,8 @@ public static <T> Read<T> read(
.setDeduplicateRecords(DEFAULT_DEDUPLICATE_RECORDS)
.setWatermarkIdleDurationThreshold(DEFAULT_WATERMARK_IDLE_DURATION_THRESHOLD)
.setAckDeadline(DEFAULT_ACK_DEADLINE)
.setNackOnTimeout(DEFAULT_NACK_ON_TIMEOUT));
.setNackOnTimeout(DEFAULT_NACK_ON_TIMEOUT)
.setEnableOpenTelemetryTracing(DEFAULT_ENABLE_OPENTELEMETRY_TRACING));
}

/**
Expand Down Expand Up @@ -706,6 +722,16 @@ public Read<T> withSessionServiceFactory(SessionServiceFactory sessionServiceFac
return this;
}

/**
* Optional, default: false. Enables OpenTelemetry distributed tracing by extracting the W3C
* trace context from incoming {@link Solace.Record#getUserProperties()} and creating a consumer
* span around downstream processing.
*/
public Read<T> withEnableOpenTelemetryTracing() {
configurationBuilder.setEnableOpenTelemetryTracing(true);
return this;
}

@AutoValue
abstract static class Configuration<T> {

Expand Down Expand Up @@ -733,9 +759,12 @@ abstract static class Configuration<T> {

abstract boolean getNackOnTimeout();

abstract boolean getEnableOpenTelemetryTracing();

public static <T> Builder<T> builder() {
Builder<T> builder =
new org.apache.beam.sdk.io.solace.AutoValue_SolaceIO_Read_Configuration.Builder<T>();
new org.apache.beam.sdk.io.solace.AutoValue_SolaceIO_Read_Configuration.Builder<T>()
.setEnableOpenTelemetryTracing(DEFAULT_ENABLE_OPENTELEMETRY_TRACING);
return builder;
}

Expand Down Expand Up @@ -767,6 +796,8 @@ abstract Builder<T> setParseFn(

abstract Builder<T> setNackOnTimeout(boolean nackOnTimeout);

abstract Builder<T> setEnableOpenTelemetryTracing(boolean enableOpenTelemetryTracing);

abstract Configuration<T> build();
}
}
Expand All @@ -793,20 +824,32 @@ public PCollection<T> expand(PBegin input) {

Coder<T> coder = inferCoder(input.getPipeline(), configuration.getTypeDescriptor());

return input.apply(
org.apache.beam.sdk.io.Read.from(
new UnboundedSolaceSource<>(
initializedQueue,
sempClientFactory,
sessionServiceFactory,
configuration.getMaxNumConnections(),
configuration.getDeduplicateRecords(),
coder,
configuration.getTimestampFn(),
configuration.getWatermarkIdleDurationThreshold(),
configuration.getParseFn(),
configuration.getAckDeadline(),
configuration.getNackOnTimeout())));
PCollection<T> output =
input.apply(
org.apache.beam.sdk.io.Read.from(
new UnboundedSolaceSource<>(
initializedQueue,
sempClientFactory,
sessionServiceFactory,
configuration.getMaxNumConnections(),
configuration.getDeduplicateRecords(),
coder,
configuration.getTimestampFn(),
configuration.getWatermarkIdleDurationThreshold(),
configuration.getParseFn(),
configuration.getAckDeadline(),
configuration.getNackOnTimeout())));

if (configuration.getEnableOpenTelemetryTracing()) {
output =
output
.apply(
"Extract OpenTelemetry context from Header",
ParDo.of(new OpenTelemetryHeaderConsumer<>()))
.setCoder(coder);
}

return output;
}

@VisibleForTesting
Expand Down Expand Up @@ -871,6 +914,118 @@ private Queue initializeQueueForTopicIfNeeded(
}
}

@VisibleForTesting
static class OpenTelemetryHeaderConsumer<T> extends DoFn<T, T> {
private transient @Nullable Tracer tracer;

@Setup
public void setup(PipelineOptions options) {
OpenTelemetry openTelemetry = options.as(SdkHarnessOptions.class).getOpenTelemetry();
tracer = openTelemetry.getTracer("SolaceIO");
}

private static final TextMapGetter<Record> HEADER_GETTER =
new TextMapGetter<Record>() {
@Override
public Iterable<String> keys(Record carrier) {
return carrier.getUserProperties().keySet();
}

@Override
public @Nullable String get(@Nullable Record carrier, String key) {
if (carrier == null) {
return null;
}
Solace.UserPropertyValue direct = carrier.getUserProperties().get(key);
if (direct != null && direct.getString() != null) {
return direct.getString();
}
for (Map.Entry<String, Solace.UserPropertyValue> entry :
carrier.getUserProperties().entrySet()) {
if (entry.getKey().equalsIgnoreCase(key)) {
return entry.getValue().getString();
}
}
return null;
}
};

@ProcessElement
public void processElement(@Element T element, OutputReceiver<T> receiver) {
Tracer currentTracer = this.tracer;
if (currentTracer == null) {
receiver.output(element);
return;
}

Context context = Context.current();
if (element instanceof Record) {
context =
W3CTraceContextPropagator.getInstance()
.extract(context, (Record) element, HEADER_GETTER);
}
Span span =
currentTracer
.spanBuilder("SolaceIO.Read")
.setSpanKind(SpanKind.CONSUMER)
.setParent(context)
.startSpan();
try (Scope scope = span.makeCurrent()) {
receiver.output(element);
} finally {
span.end();
}
}
}

@VisibleForTesting
static class OpenTelemetryHeaderPropagator extends DoFn<Record, Record> {
private transient @Nullable Tracer tracer;

@Setup
public void setup(PipelineOptions options) {
OpenTelemetry openTelemetry = options.as(SdkHarnessOptions.class).getOpenTelemetry();
tracer = openTelemetry.getTracer("SolaceIO");
}

private static final TextMapSetter<Record> HEADER_SETTER =
new TextMapSetter<Record>() {
@Override
public void set(@Nullable Record carrier, String key, String value) {
if (carrier != null && value != null) {
carrier.getUserProperties().put(key, Solace.UserPropertyValue.of(value));
}
}
};

@ProcessElement
public void processElement(@Element Record element, OutputReceiver<Record> receiver) {
Tracer currentTracer = this.tracer;
if (currentTracer == null) {
receiver.output(element);
return;
}

Span span =
currentTracer
.spanBuilder("SolaceIO.Write")
.setSpanKind(SpanKind.PRODUCER)
.setParent(Context.current())
.startSpan();
try (Scope scope = span.makeCurrent()) {
Record updatedElement =
element.toBuilder()
.setUserProperties(new HashMap<>(element.getUserProperties()))
.build();
W3CTraceContextPropagator.getInstance()
.inject(Context.current(), updatedElement, HEADER_SETTER);
receiver.output(updatedElement);
} finally {
span.end();
}
}
}

public enum SubmissionMode {
HIGHER_THROUGHPUT,
LOWER_LATENCY,
Expand Down Expand Up @@ -1061,6 +1216,15 @@ public Write<T> withErrorHandler(ErrorHandler<BadRecord, ?> errorHandler) {
return toBuilder().setErrorHandler(errorHandler).build();
}

/**
* Optional, default: false. Enables OpenTelemetry distributed tracing by creating a producer
* span and injecting the W3C trace context into outgoing {@link
* Solace.Record#getUserProperties()}.
*/
public Write<T> withEnableOpenTelemetryTracing() {
return toBuilder().setEnableOpenTelemetryTracing(true).build();
}

abstract int getNumShards();

abstract int getNumberOfClientsPerWorker();
Expand All @@ -1081,14 +1245,17 @@ public Write<T> withErrorHandler(ErrorHandler<BadRecord, ?> errorHandler) {

abstract @Nullable ErrorHandler<BadRecord, ?> getErrorHandler();

abstract boolean getEnableOpenTelemetryTracing();

static <T> Builder<T> builder() {
return new AutoValue_SolaceIO_Write.Builder<T>()
.setDeliveryMode(DEFAULT_WRITER_DELIVERY_MODE)
.setNumShards(DEFAULT_WRITER_NUM_SHARDS)
.setNumberOfClientsPerWorker(DEFAULT_WRITER_CLIENTS_PER_WORKER)
.setPublishLatencyMetrics(DEFAULT_WRITER_PUBLISH_LATENCY_METRICS)
.setDispatchMode(DEFAULT_WRITER_SUBMISSION_MODE)
.setWriterType(DEFAULT_WRITER_TYPE);
.setWriterType(DEFAULT_WRITER_TYPE)
.setEnableOpenTelemetryTracing(DEFAULT_ENABLE_OPENTELEMETRY_TRACING);
}

abstract Builder<T> toBuilder();
Expand All @@ -1115,6 +1282,8 @@ abstract static class Builder<T> {

abstract Builder<T> setErrorHandler(ErrorHandler<BadRecord, ?> errorHandler);

abstract Builder<T> setEnableOpenTelemetryTracing(boolean enableOpenTelemetryTracing);

abstract Write<T> build();
}

Expand Down Expand Up @@ -1146,6 +1315,12 @@ public SolaceOutput expand(PCollection<T> input) {
MapElements.into(TypeDescriptor.of(Solace.Record.class))
.via(checkNotNull(getFormatFunction())));

if (getEnableOpenTelemetryTracing()) {
records =
records.apply(
"Propagate OpenTelemetry Tracing", ParDo.of(new OpenTelemetryHeaderPropagator()));
}

PCollection<Solace.Record> withGlobalWindow =
records.apply("Global window", Window.into(new GlobalWindows()));

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -435,6 +435,8 @@ public static Builder builder() {
.setUserProperties(Collections.emptyMap());
}

public abstract Builder toBuilder();

@AutoValue.Builder
public abstract static class Builder {
public abstract Builder setMessageId(String messageId);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,8 @@
import com.solacesystems.jcsmp.Destination;
import com.solacesystems.jcsmp.JCSMPException;
import java.time.Instant;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import org.apache.beam.sdk.io.solace.broker.MessageProducer;
import org.apache.beam.sdk.io.solace.broker.PublishResultHandler;
Expand All @@ -29,6 +31,19 @@
import org.apache.beam.sdk.transforms.SerializableFunction;

public abstract class MockProducer implements MessageProducer {
private static final List<Record> PUBLISHED_RECORDS =
Collections.synchronizedList(new ArrayList<>());

public static List<Record> getPublishedRecords() {
synchronized (PUBLISHED_RECORDS) {
return new ArrayList<>(PUBLISHED_RECORDS);
}
}

public static void clearPublishedRecords() {
PUBLISHED_RECORDS.clear();
}

final PublishResultHandler handler;

public MockProducer(PublishResultHandler handler) {
Expand Down Expand Up @@ -67,6 +82,7 @@ public void publishSingleMessage(
Destination topicOrQueue,
boolean useCorrelationKeyLatency,
DeliveryMode deliveryMode) {
PUBLISHED_RECORDS.add(msg);
if (useCorrelationKeyLatency) {
handler.responseReceivedEx(
Solace.PublishResult.builder()
Expand Down
Loading
Loading