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]

Reply via email to