stankiewicz commented on PR #40187:
URL: https://github.com/apache/beam/pull/40187#issuecomment-5832229972
The implementation is incomplete - In its current state, the PR introduces a
configuration toggle (withEnableOpenTelemetryTracing()) and passes it down to
UnboundedSolaceReader, but no actual tracing instrumentation is performed. The
reader merely logs a debug message (LOG.debug(...)), without creating
OpenTelemetry spans, propagating W3C trace context, or supporting write
pipelines.
key flaws:
1. No Actual OpenTelemetry Instrumentation -No spans are created or
recorded, no trace parent is extracted, and no OpenTelemetry API/SDK
dependencies are imported. The feature is essentially a no-op that produces log.
2. Tracing Inside UnboundedSolaceReader Instead of a ParDo. KafkaIO
and PubsubIO handle distributed tracing via a dedicated ParDo transform
applied to the PCollection in expand(). This ensures proper span lifecycle
scoping across element processing.
3. Completely Missing Write Support (SolaceIO.Write) , FR was for both Read
and Write
4. Ignoring Solace Message User Properties (PR #40108) - Distributed tracing
in messaging systems relies on W3C Trace Context headers (traceparent,
tracestate). [PR #40108](https://github.com/apache/beam/pull/40108) was merged
to expose Solace JCSMP user properties as a Map<String, String> on
Solace.Record. Without reading or writing these properties, W3C trace
propagation cannot function.
5. Unrelated Build System Changes
Changes requested:
1. Add missing OpenTelemetry dependencies in sdks/java/io/solace/build.gradle
- implementation platform(library.java.opentelemetry_bom)
- implementation library.java.opentelemetry_api
- implementation library.java.opentelemetry_context
- testImplementation library.java.opentelemetry_sdk
2. Read-Side Trace Extraction
A dedicated DoFn (e.g., OpenTelemetryHeaderConsumer) applied in
SolaceIO.Read.expand().
OpenTelemetry instance injection via PipelineOptions:
```
openTelemetry = options.as(SdkHarnessOptions.class).getOpenTelemetry();
tracer = openTelemetry.getTracer("SolaceIO");
```
Extraction of W3C trace context from `record.getProperties()` using
`W3CTraceContextPropagator.getInstance().extract(...)` and
`TextMapGetter<Solace.Record>`.
Creating a consumer span ("SolaceIO.Read") wrapping
`receiver.output(record)`.
3. Write-Side Trace Injection
- withEnableOpenTelemetryTracing() on SolaceIO.Write.
- A dedicated DoFn (e.g., OpenTelemetryHeaderPropagator) applied in
`SolaceIO.Write.expand()`.
- Injection of the current trace context into Solace.Record user properties
using `W3CTraceContextPropagator.getInstance().inject(...)` and `TextMapSetter`.
Align the implementation with the pattern used in with KafkaIO.
4. Revert
- Revert changes to BeamModulePlugin.groovy and keep the PR focused on
SolaceIO.
- Revert the changes in UnboundedSolaceReader and UnboundedSolaceSource to
prevent mixing transport polling with telemetry context propagation.
5. Tests
Add unit tests with InMemorySpanExporter to verify span creation and context
propagation.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]