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.
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
}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")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");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");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");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")#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");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?;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'
)# 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.
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())
}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))
raiseawait 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();
}
});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();
}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;
}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()
}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();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();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
endrequire 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
endMetrics
Instrument your messaging clients with OpenTelemetry metrics to track throughput, latency, and error rates.
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
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 |
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?