KubeMQ
LearnGuides

OpenTelemetry Integration

Add distributed tracing and metrics to KubeMQ messaging with OpenTelemetry.

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.

Tracing is how you observe message flow under load. See Scaling & Flow in Fundamentals for the concepts behind throughput, backpressure, and consumer groups that these traces will surface.

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

Setup

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

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
}
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")
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");
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");
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");
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")
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");
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?;
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'
)
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")

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 to convert these into spans, or attach your own handler.

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.

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())
}
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
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();
  }
});
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();
}
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;
}
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()
}
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();
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();
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
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

Metrics

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

MetricTypeDescription
messaging.publish.durationHistogramTime to publish a message (ms)
messaging.process.durationHistogramTime to process a received message (ms)
messaging.publish.messagesCounterTotal messages published
messaging.receive.messagesCounterTotal messages received
messaging.publish.errorsCounterFailed publish attempts

Attribute Conventions

Follow OpenTelemetry semantic conventions for messaging attributes:

AttributeExampleDescription
messaging.systemkubemqMessaging system identifier
messaging.destinationorder-eventsChannel name
messaging.operationpublish / receiveOperation type
messaging.message.iduuidMessage identifier

These conventions align with the OpenTelemetry Semantic Conventions for Messaging. Following them ensures compatibility with observability platforms like Jaeger, Grafana Tempo, and Datadog.

Next Steps

Was this page helpful?

On this page