# OpenTelemetry Integration (/learn/guides/opentelemetry)



## Overview [#overview]

KubeMQ supports OpenTelemetry for distributed tracing and metrics. Instrumenting your producers and consumers gives you end-to-end visibility across messaging boundaries — regardless of which pattern (Events, Events Store, Queues, or RPC) you use.

<Callout type="info">
  Tracing is how you observe message flow under load. See [Scaling & Flow](/learn/concepts/scaling-and-flow) in Fundamentals for the concepts behind throughput, backpressure, and consumer groups that these traces will surface.
</Callout>

<Mermaid
  chart="graph LR
    P[&#x22;Producer&#x22;]
    K[&#x22;KubeMQ&#x22;]
    C[&#x22;Consumer&#x22;]
    OT[&#x22;OTel Collector&#x22;]
    J[&#x22;Jaeger / Tempo / etc.&#x22;]

    P -- &#x22;traced send&#x22; --> K
    K -- &#x22;traced deliver&#x22; --> C
    P -. &#x22;spans&#x22; .-> OT
    C -. &#x22;spans&#x22; .-> OT
    OT -- &#x22;export&#x22; --> J

    class P,C client
    class K broker
    class OT,J external"
/>

*Producer and consumer spans flow through the OTel Collector to a tracing backend, giving end-to-end visibility across the KubeMQ message path.*

## Setup [#setup]

Configure an OpenTelemetry tracer and meter provider in your application before creating KubeMQ clients.

<Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
  <Tab value="Go">
    ```go title="otel_setup.go"
    package main

    import (
        "context"
        "log"

        "go.opentelemetry.io/otel"
        "go.opentelemetry.io/otel/exporters/otlp/otlptrace/otlptracegrpc"
        "go.opentelemetry.io/otel/sdk/resource"
        sdktrace "go.opentelemetry.io/otel/sdk/trace"
        semconv "go.opentelemetry.io/otel/semconv/v1.24.0"
    )

    func initTracer(ctx context.Context) (*sdktrace.TracerProvider, error) {
        exporter, err := otlptracegrpc.New(ctx,
            otlptracegrpc.WithEndpoint("localhost:4317"),
            otlptracegrpc.WithInsecure(),
        )
        if err != nil {
            return nil, err
        }

        tp := sdktrace.NewTracerProvider(
            sdktrace.WithBatcher(exporter),
            sdktrace.WithResource(resource.NewWithAttributes(
                semconv.SchemaURL,
                semconv.ServiceName("order-service"),
            )),
        )
        otel.SetTracerProvider(tp)
        return tp, nil
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python title="otel_setup.py"
    from opentelemetry import trace
    from opentelemetry.sdk.trace import TracerProvider
    from opentelemetry.sdk.trace.export import BatchSpanProcessor
    from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
    from opentelemetry.sdk.resources import Resource

    resource = Resource.create({"service.name": "order-service"})
    provider = TracerProvider(resource=resource)

    exporter = OTLPSpanExporter(endpoint="localhost:4317", insecure=True)
    provider.add_span_processor(BatchSpanProcessor(exporter))

    trace.set_tracer_provider(provider)
    tracer = trace.get_tracer("order-service")
    ```
  </Tab>

  <Tab value="Node.js">
    ```typescript title="otel_setup.ts"
    import { NodeTracerProvider } from "@opentelemetry/sdk-trace-node";
    import { BatchSpanProcessor } from "@opentelemetry/sdk-trace-base";
    import { OTLPTraceExporter } from "@opentelemetry/exporter-trace-otlp-grpc";
    import { Resource } from "@opentelemetry/resources";
    import { ATTR_SERVICE_NAME } from "@opentelemetry/semantic-conventions";

    const provider = new NodeTracerProvider({
      resource: new Resource({ [ATTR_SERVICE_NAME]: "order-service" }),
    });

    provider.addSpanProcessor(
      new BatchSpanProcessor(
        new OTLPTraceExporter({ url: "http://localhost:4317" })
      )
    );
    provider.register();

    const tracer = provider.getTracer("order-service");
    ```
  </Tab>

  <Tab value="Java">
    ```java title="OtelSetup.java"
    import io.opentelemetry.api.OpenTelemetry;
    import io.opentelemetry.api.trace.Tracer;
    import io.opentelemetry.exporter.otlp.trace.OtlpGrpcSpanExporter;
    import io.opentelemetry.sdk.OpenTelemetrySdk;
    import io.opentelemetry.sdk.resources.Resource;
    import io.opentelemetry.sdk.trace.SdkTracerProvider;
    import io.opentelemetry.sdk.trace.export.BatchSpanProcessor;
    import io.opentelemetry.semconv.ResourceAttributes;

    Resource resource = Resource.getDefault()
        .merge(Resource.create(
            io.opentelemetry.api.common.Attributes.of(
                ResourceAttributes.SERVICE_NAME, "order-service")));

    OtlpGrpcSpanExporter exporter = OtlpGrpcSpanExporter.builder()
        .setEndpoint("http://localhost:4317")
        .build();

    SdkTracerProvider tracerProvider = SdkTracerProvider.builder()
        .addSpanProcessor(BatchSpanProcessor.builder(exporter).build())
        .setResource(resource)
        .build();

    OpenTelemetry otel = OpenTelemetrySdk.builder()
        .setTracerProvider(tracerProvider)
        .build();

    Tracer tracer = otel.getTracer("order-service");
    ```
  </Tab>

  <Tab value="C#">
    ```csharp title="OtelSetup.cs"
    using OpenTelemetry;
    using OpenTelemetry.Resources;
    using OpenTelemetry.Trace;

    using var tracerProvider = Sdk.CreateTracerProviderBuilder()
        .SetResourceBuilder(ResourceBuilder.CreateDefault()
            .AddService("order-service"))
        .AddOtlpExporter(opts =>
        {
            opts.Endpoint = new Uri("http://localhost:4317");
        })
        .Build();

    var tracer = tracerProvider.GetTracer("order-service");
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin title="OtelSetup.kt"
    import io.opentelemetry.api.OpenTelemetry
    import io.opentelemetry.api.trace.Tracer
    import io.opentelemetry.exporter.otlp.trace.OtlpGrpcSpanExporter
    import io.opentelemetry.sdk.OpenTelemetrySdk
    import io.opentelemetry.sdk.resources.Resource
    import io.opentelemetry.sdk.trace.SdkTracerProvider
    import io.opentelemetry.sdk.trace.export.BatchSpanProcessor
    import io.opentelemetry.semconv.ResourceAttributes

    val resource = Resource.getDefault()
        .merge(Resource.create(
            io.opentelemetry.api.common.Attributes.of(
                ResourceAttributes.SERVICE_NAME, "order-service")))

    val exporter = OtlpGrpcSpanExporter.builder()
        .setEndpoint("http://localhost:4317")
        .build()

    val tracerProvider = SdkTracerProvider.builder()
        .addSpanProcessor(BatchSpanProcessor.builder(exporter).build())
        .setResource(resource)
        .build()

    val otel: OpenTelemetry = OpenTelemetrySdk.builder()
        .setTracerProvider(tracerProvider)
        .build()

    val tracer: Tracer = otel.getTracer("order-service")
    ```
  </Tab>

  <Tab value="C++">
    ```cpp title="otel_setup.cpp"
    #include <opentelemetry/sdk/trace/tracer_provider.h>
    #include <opentelemetry/exporters/otlp/otlp_grpc_exporter.h>
    #include <opentelemetry/sdk/trace/batch_span_processor.h>
    #include <opentelemetry/sdk/resource/resource.h>
    #include <opentelemetry/trace/provider.h>

    namespace trace_sdk = opentelemetry::sdk::trace;
    namespace otlp = opentelemetry::exporter::otlp;
    namespace resource = opentelemetry::sdk::resource;

    auto exporter = std::make_unique<otlp::OtlpGrpcExporter>(
        otlp::OtlpGrpcExporterOptions{"localhost:4317"});

    auto processor = std::make_unique<trace_sdk::BatchSpanProcessor>(
        std::move(exporter));

    auto provider = std::make_shared<trace_sdk::TracerProvider>(
        std::move(processor),
        resource::Resource::Create({{"service.name", "order-service"}}));

    opentelemetry::trace::Provider::SetTracerProvider(provider);
    auto tracer = provider->GetTracer("order-service");
    ```
  </Tab>

  <Tab value="Rust">
    ```rust title="otel_setup.rs"
    use kubemq::prelude::*;
    use opentelemetry::global;
    use opentelemetry_otlp::WithExportConfig;
    use opentelemetry_sdk::{trace::TracerProvider, Resource};
    use opentelemetry::KeyValue;

    // Register a global TracerProvider before building the client. Once a
    // provider is registered globally, the KubeMQ Rust client automatically
    // instruments its gRPC calls — no per-call wiring required.
    let exporter = opentelemetry_otlp::SpanExporter::builder()
        .with_tonic()
        .with_endpoint("http://localhost:4317")
        .build()?;

    let provider = TracerProvider::builder()
        .with_batch_exporter(exporter, opentelemetry_sdk::runtime::Tokio)
        .with_resource(Resource::new(vec![
            KeyValue::new("service.name", "order-service"),
        ]))
        .build();

    global::set_tracer_provider(provider);

    let client = KubemqClient::builder()
        .host("localhost")
        .port(50000)
        .build()
        .await?;
    ```
  </Tab>

  <Tab value="Ruby">
    ```ruby title="otel_setup.rb"
    require 'kubemq'
    require 'opentelemetry/sdk'
    require 'opentelemetry/exporter/otlp'

    # Configure the global OpenTelemetry SDK before creating the client. With a
    # provider configured, the KubeMQ Ruby client's gRPC calls are instrumented
    # through the standard OpenTelemetry gRPC instrumentation.
    OpenTelemetry::SDK.configure do |c|
      c.service_name = 'order-service'
      c.use_all
    end

    tracer = OpenTelemetry.tracer_provider.tracer('order-service')

    client = KubeMQ::PubSubClient.new(
      address: 'localhost:50000',
      client_id: 'order-service'
    )
    ```
  </Tab>

  <Tab value="Elixir">
    ```elixir title="otel_setup.exs"
    # The KubeMQ Elixir client emits :telemetry events for every operation
    # rather than OpenTelemetry spans directly. Bridge those telemetry events
    # into OpenTelemetry with the opentelemetry_telemetry library, or attach a
    # handler to forward measurements to your tracer.
    :telemetry.attach_many(
      "kubemq-otel-bridge",
      [
        [:kubemq, :client, :send_event, :start],
        [:kubemq, :client, :send_event, :stop],
        [:kubemq, :client, :send_event, :exception]
      ],
      fn event, measurements, metadata, _config ->
        [_, _, action, phase] = event
        # Forward to your OpenTelemetry tracer here, e.g. via OpenTelemetry.Tracer.
        IO.puts("[otel] #{action}:#{phase} #{inspect(metadata[:channel])}")
      end,
      nil
    )

    {:ok, client} =
      KubeMQ.Client.start_link(address: "localhost:50000", client_id: "order-service")
    ```
  </Tab>
</Tabs>

<Callout type="info">
  The Elixir client exposes observability through Erlang/Elixir `:telemetry` events (`[:kubemq, :client, :send_event, :start | :stop | :exception]`, and equivalents for commands, queries, and queues) rather than emitting OpenTelemetry spans directly. Use [`opentelemetry_telemetry`](https://hex.pm/packages/opentelemetry_telemetry) to convert these into spans, or attach your own handler.
</Callout>

## Tracing Messaging Operations [#tracing-messaging-operations]

Wrap your publish and subscribe operations in spans to trace message flow across services. The pattern is the same for Events, Events Store, Queues, and RPC.

<Tabs groupId="language" items="['Go', 'Python', 'Node.js', 'Java', 'C#', 'Kotlin', 'C++', 'Rust', 'Ruby', 'Elixir']">
  <Tab value="Go">
    ```go title="traced_publish.go"
    tracer := otel.Tracer("order-service")

    ctx, span := tracer.Start(ctx, "publish-order",
        trace.WithAttributes(
            attribute.String("messaging.system", "kubemq"),
            attribute.String("messaging.destination", "order-events"),
            attribute.String("messaging.operation", "publish"),
        ),
    )
    defer span.End()

    err := client.SendEvent(ctx, kubemq.NewEvent().
        SetChannel("order-events").
        SetBody([]byte(`{"orderId":"ORD-500"}`)),
    )
    if err != nil {
        span.RecordError(err)
        span.SetStatus(codes.Error, err.Error())
    }
    ```
  </Tab>

  <Tab value="Python">
    ```python title="traced_publish.py"
    with tracer.start_as_current_span(
        "publish-order",
        attributes={
            "messaging.system": "kubemq",
            "messaging.destination": "order-events",
            "messaging.operation": "publish",
        },
    ) as span:
        try:
            client.send_event(
                EventMessage(
                    channel="order-events",
                    body=b'{"orderId":"ORD-500"}',
                )
            )
        except Exception as e:
            span.record_exception(e)
            span.set_status(StatusCode.ERROR, str(e))
            raise
    ```
  </Tab>

  <Tab value="Node.js">
    ```typescript title="traced_publish.ts"
    await tracer.startActiveSpan("publish-order", {
      attributes: {
        "messaging.system": "kubemq",
        "messaging.destination": "order-events",
        "messaging.operation": "publish",
      },
    }, async (span) => {
      try {
        await client.sendEvent({
          channel: "order-events",
          body: Buffer.from('{"orderId":"ORD-500"}'),
        });
      } catch (err) {
        span.recordException(err);
        span.setStatus({ code: SpanStatusCode.ERROR, message: err.message });
        throw err;
      } finally {
        span.end();
      }
    });
    ```
  </Tab>

  <Tab value="Java">
    ```java title="TracedPublish.java"
    Span span = tracer.spanBuilder("publish-order")
        .setAttribute("messaging.system", "kubemq")
        .setAttribute("messaging.destination", "order-events")
        .setAttribute("messaging.operation", "publish")
        .startSpan();

    try (Scope scope = span.makeCurrent()) {
        client.sendEventsMessage(EventMessage.builder()
            .channel("order-events")
            .body("{\"orderId\":\"ORD-500\"}".getBytes())
            .build());
    } catch (Exception e) {
        span.recordException(e);
        span.setStatus(StatusCode.ERROR, e.getMessage());
        throw e;
    } finally {
        span.end();
    }
    ```
  </Tab>

  <Tab value="C#">
    ```csharp title="TracedPublish.cs"
    using var span = tracer.StartActiveSpan("publish-order");
    span.SetAttribute("messaging.system", "kubemq");
    span.SetAttribute("messaging.destination", "order-events");
    span.SetAttribute("messaging.operation", "publish");

    try
    {
        await client.SendEventAsync(new EventMessage
        {
            Channel = "order-events",
            Body = Encoding.UTF8.GetBytes("{\"orderId\":\"ORD-500\"}"),
        });
    }
    catch (Exception ex)
    {
        span.RecordException(ex);
        span.SetStatus(Status.Error.WithDescription(ex.Message));
        throw;
    }
    ```
  </Tab>

  <Tab value="Kotlin">
    ```kotlin title="TracedPublish.kt"
    val span = tracer.spanBuilder("publish-order")
        .setAttribute("messaging.system", "kubemq")
        .setAttribute("messaging.destination", "order-events")
        .setAttribute("messaging.operation", "publish")
        .startSpan()

    try {
        span.makeCurrent().use {
            client.sendEvent(EventMessage(
                channel = "order-events",
                body = """{"orderId":"ORD-500"}""".toByteArray(),
            ))
        }
    } catch (e: Exception) {
        span.recordException(e)
        span.setStatus(StatusCode.ERROR, e.message ?: "unknown error")
        throw e
    } finally {
        span.end()
    }
    ```
  </Tab>

  <Tab value="C++">
    ```cpp title="traced_publish.cpp"
    auto span = tracer->StartSpan("publish-order", {
        {"messaging.system", "kubemq"},
        {"messaging.destination", "order-events"},
        {"messaging.operation", "publish"},
    });
    auto scope = tracer->WithActiveSpan(span);

    try {
        kubemq::EventMessage event;
        event.channel = "order-events";
        event.body = R"({"orderId":"ORD-500"})";
        client.sendEvent(event);
    } catch (const std::exception& e) {
        span->AddEvent("exception", {{"exception.message", e.what()}});
        span->SetStatus(opentelemetry::trace::StatusCode::kError, e.what());
        throw;
    }
    span->End();
    ```
  </Tab>

  <Tab value="Rust">
    ```rust title="traced_publish.rs"
    use kubemq::EventBuilder;
    use opentelemetry::trace::{Span, Status, Tracer};
    use opentelemetry::{global, KeyValue};

    let tracer = global::tracer("order-service");
    let mut span = tracer.start("publish-order");
    span.set_attribute(KeyValue::new("messaging.system", "kubemq"));
    span.set_attribute(KeyValue::new("messaging.destination", "order-events"));
    span.set_attribute(KeyValue::new("messaging.operation", "publish"));

    let event = EventBuilder::new()
        .channel("order-events")
        .body(br#"{"orderId":"ORD-500"}"#.to_vec())
        .build();

    match client.send_event(event).await {
        Ok(_) => {}
        Err(e) => {
            span.set_status(Status::error(e.to_string()));
            span.record_error(&e);
        }
    }
    span.end();
    ```
  </Tab>

  <Tab value="Ruby">
    ```ruby title="traced_publish.rb"
    tracer.in_span(
      'publish-order',
      attributes: {
        'messaging.system' => 'kubemq',
        'messaging.destination' => 'order-events',
        'messaging.operation' => 'publish'
      }
    ) do |span|
      begin
        msg = KubeMQ::PubSub::EventMessage.new(
          channel: 'order-events',
          body: '{"orderId":"ORD-500"}'
        )
        client.send_event(msg)
      rescue KubeMQ::Error => e
        span.record_exception(e)
        span.status = OpenTelemetry::Trace::Status.error(e.message)
        raise
      end
    end
    ```
  </Tab>

  <Tab value="Elixir">
    ```elixir title="traced_publish.exs"
    require OpenTelemetry.Tracer, as: Tracer

    Tracer.with_span "publish-order" do
      Tracer.set_attributes([
        {"messaging.system", "kubemq"},
        {"messaging.destination", "order-events"},
        {"messaging.operation", "publish"}
      ])

      event =
        KubeMQ.Event.new(channel: "order-events", body: ~s({"orderId":"ORD-500"}))

      case KubeMQ.Client.send_event(client, event) do
        :ok ->
          :ok

        {:error, err} ->
          Tracer.set_status(OpenTelemetry.status(:error, err.message))
      end
    end
    ```
  </Tab>
</Tabs>

## Metrics [#metrics]

Instrument your messaging clients with OpenTelemetry metrics to track throughput, latency, and error rates.

### Recommended Metrics [#recommended-metrics]

| Metric                       | Type      | Description                             |
| ---------------------------- | --------- | --------------------------------------- |
| `messaging.publish.duration` | Histogram | Time to publish a message (ms)          |
| `messaging.process.duration` | Histogram | Time to process a received message (ms) |
| `messaging.publish.messages` | Counter   | Total messages published                |
| `messaging.receive.messages` | Counter   | Total messages received                 |
| `messaging.publish.errors`   | Counter   | Failed publish attempts                 |

### Attribute Conventions [#attribute-conventions]

Follow OpenTelemetry semantic conventions for messaging attributes:

| Attribute               | Example               | Description                 |
| ----------------------- | --------------------- | --------------------------- |
| `messaging.system`      | `kubemq`              | Messaging system identifier |
| `messaging.destination` | `order-events`        | Channel name                |
| `messaging.operation`   | `publish` / `receive` | Operation type              |
| `messaging.message.id`  | `uuid`                | Message identifier          |

<Callout type="info">
  These conventions align with the [OpenTelemetry Semantic Conventions for Messaging](https://opentelemetry.io/docs/specs/semconv/messaging/). Following them ensures compatibility with observability platforms like Jaeger, Grafana Tempo, and Datadog.
</Callout>

## Next Steps [#next-steps]

<Cards>
  <Card title="Error Handling" href="/learn/guides/error-handling" description="Use traced errors alongside retry logic." />

  <Card title="Production Checklist" href="/learn/guides/production-checklist" description="Verify observability configuration before going live." />
</Cards>
