KubeMQ
LearnEventsScenarios

Live Dashboard Data Feed

Stream real-time metrics to a dashboard using KubeMQ Events.

Architecture

Microservices publish operational metrics to KubeMQ. A dashboard aggregator subscribes with a wildcard and feeds real-time data to a frontend display.

Each service publishes to a hierarchical metrics.* channel; the aggregator subscribes once with the metrics.> wildcard and pushes live updates to the UI.

Implementation

Metrics Publisher

Each service publishes metrics to hierarchical channels for flexible subscription.

metrics_publisher.go
package main

import (
    "context"
    "fmt"
    "log"
    "math/rand"
    "time"

    "github.com/kubemq-io/kubemq-go/v2"
)

func main() {
    ctx := context.Background()
    client, err := kubemq.NewClient(ctx,
        kubemq.WithAddress("localhost", 50000),
    )
    if err != nil {
        log.Fatal(err)
    }
    defer client.Close()

    ticker := time.NewTicker(2 * time.Second)
    defer ticker.Stop()

    for range ticker.C {
        metrics := []struct {
            channel string
            value   float64
        }{
            {"metrics.orders.count", float64(rand.Intn(100))},
            {"metrics.orders.revenue", float64(rand.Intn(10000))},
            {"metrics.inventory.stock_level", float64(rand.Intn(500))},
        }

        for _, m := range metrics {
            body := fmt.Sprintf(`{"value":%.2f,"timestamp":%d}`, m.value, time.Now().UnixMilli())
            err := client.SendEvent(ctx, kubemq.NewEvent().
                SetChannel(m.channel).
                SetBody([]byte(body)),
            )
            if err != nil {
                log.Printf("Failed to publish %s: %v", m.channel, err)
            }
        }
    }
}
metrics_publisher.py
import json
import time
import random
from kubemq.pubsub import Client as PubSubClient
from kubemq.pubsub import EventMessage

client = PubSubClient(address="localhost:50000")

while True:
    metrics = [
        ("metrics.orders.count", random.randint(0, 100)),
        ("metrics.orders.revenue", random.randint(0, 10000)),
        ("metrics.inventory.stock_level", random.randint(0, 500)),
    ]

    for channel, value in metrics:
        body = json.dumps({"value": value, "timestamp": int(time.time() * 1000)})
        client.send_event(
            EventMessage(channel=channel, body=body.encode("utf-8"))
        )

    time.sleep(2)
metrics_publisher.js
const { KubeMQClient } = require("kubemq-js");

const client = new KubeMQClient({ address: "localhost:50000" });

setInterval(async () => {
  const metrics = [
    { channel: "metrics.orders.count", value: Math.floor(Math.random() * 100) },
    { channel: "metrics.orders.revenue", value: Math.floor(Math.random() * 10000) },
    { channel: "metrics.inventory.stock_level", value: Math.floor(Math.random() * 500) },
  ];

  for (const m of metrics) {
    await client.sendEvent({
      channel: m.channel,
      body: Buffer.from(JSON.stringify({ value: m.value, timestamp: Date.now() })),
    });
  }
}, 2000);
MetricsPublisher.java
PubSubClient client = PubSubClient.builder()
    .address("localhost:50000")
    .clientId("metrics-publisher")
    .build();

Random rng = new Random();
while (true) {
    String[][] metrics = {
        {"metrics.orders.count", String.valueOf(rng.nextInt(100))},
        {"metrics.orders.revenue", String.valueOf(rng.nextInt(10000))},
        {"metrics.inventory.stock_level", String.valueOf(rng.nextInt(500))},
    };

    for (String[] m : metrics) {
        String body = String.format(
            "{\"value\":%s,\"timestamp\":%d}", m[1], System.currentTimeMillis());
        client.sendEventsMessage(EventMessage.builder()
            .channel(m[0]).body(body.getBytes()).build());
    }
    Thread.sleep(2000);
}
MetricsPublisher.cs
await using var client = new KubeMQClient(new KubeMQClientOptions());
await client.ConnectAsync();

var rng = new Random();
while (true)
{
    var metrics = new[]
    {
        ("metrics.orders.count", rng.Next(100)),
        ("metrics.orders.revenue", rng.Next(10000)),
        ("metrics.inventory.stock_level", rng.Next(500)),
    };

    foreach (var (channel, value) in metrics)
    {
        await client.SendEventAsync(new EventMessage
        {
            Channel = channel,
            Body = Encoding.UTF8.GetBytes(
                $"{{\"value\":{value},\"timestamp\":{DateTimeOffset.UtcNow.ToUnixTimeMilliseconds()}}}"),
        });
    }
    await Task.Delay(2000);
}
MetricsPublisher.kt
val client = PubSubClient("localhost:50000")
val rng = java.util.Random()

while (true) {
    val metrics = listOf(
        "metrics.orders.count" to rng.nextInt(100),
        "metrics.orders.revenue" to rng.nextInt(10000),
        "metrics.inventory.stock_level" to rng.nextInt(500),
    )

    for ((channel, value) in metrics) {
        val body = """{"value":$value,"timestamp":${System.currentTimeMillis()}}"""
        client.sendEvent(EventMessage(channel = channel, body = body.toByteArray()))
    }
    Thread.sleep(2000)
}
metrics_publisher.cpp
auto client = kubemq::PubSubClient("localhost:50000");

while (true) {
    std::vector<std::pair<std::string, int>> metrics = {
        {"metrics.orders.count", rand() % 100},
        {"metrics.orders.revenue", rand() % 10000},
        {"metrics.inventory.stock_level", rand() % 500},
    };

    for (const auto& [channel, value] : metrics) {
        kubemq::EventMessage event;
        event.channel = channel;
        event.body = "{\"value\":" + std::to_string(value) + "}";
        client.sendEvent(event);
    }
    std::this_thread::sleep_for(std::chrono::seconds(2));
}
metrics_publisher.rs
use kubemq::prelude::*;
use kubemq::EventBuilder;
use std::time::Duration;

#[tokio::main]
async fn main() -> kubemq::Result<()> {
    let client = KubemqClient::builder()
        .host("localhost")
        .port(50000)
        .build()
        .await?;

    let mut ticker = tokio::time::interval(Duration::from_secs(2));
    loop {
        ticker.tick().await;

        let metrics = [
            ("metrics.orders.count", rand::random::<u8>() as i32),
            ("metrics.orders.revenue", rand::random::<u16>() as i32),
            ("metrics.inventory.stock_level", rand::random::<u16>() as i32 % 500),
        ];

        for (channel, value) in metrics {
            let ts = std::time::SystemTime::now()
                .duration_since(std::time::UNIX_EPOCH)
                .unwrap()
                .as_millis();
            let body = format!(r#"{{"value":{value},"timestamp":{ts}}}"#);
            let event = EventBuilder::new()
                .channel(channel)
                .body(body.into_bytes())
                .build();
            client.send_event(event).await?;
        }
    }
}
metrics_publisher.rb
require 'kubemq'
require 'json'

client = KubeMQ::PubSubClient.new(address: 'localhost:50000', client_id: 'metrics-publisher')

loop do
  metrics = [
    ['metrics.orders.count', rand(100)],
    ['metrics.orders.revenue', rand(10_000)],
    ['metrics.inventory.stock_level', rand(500)],
  ]

  metrics.each do |channel, value|
    body = JSON.generate(value: value, timestamp: (Time.now.to_f * 1000).to_i)
    client.send_event(KubeMQ::PubSub::EventMessage.new(channel: channel, body: body))
  end

  sleep 2
end
metrics_publisher.exs
{:ok, client} =
  KubeMQ.Client.start_link(address: "localhost:50000", client_id: "metrics-publisher")

metrics_loop = fn loop ->
  metrics = [
    {"metrics.orders.count", :rand.uniform(100)},
    {"metrics.orders.revenue", :rand.uniform(10_000)},
    {"metrics.inventory.stock_level", :rand.uniform(500)}
  ]

  for {channel, value} <- metrics do
    body = Jason.encode!(%{value: value, timestamp: System.system_time(:millisecond)})
    event = KubeMQ.Event.new(channel: channel, body: body)
    :ok = KubeMQ.Client.send_event(client, event)
  end

  Process.sleep(2000)
  loop.(loop)
end

metrics_loop.(metrics_loop)

Dashboard Aggregator

Subscribe with a wildcard to capture all metrics and aggregate them for the UI.

dashboard_aggregator.go
type DashboardState struct {
    mu      sync.RWMutex
    metrics map[string]float64
}

state := &DashboardState{metrics: make(map[string]float64)}

sub, err := client.SubscribeToEvents(ctx, "metrics.>", "",
    kubemq.WithOnEvent(func(event *kubemq.Event) {
        var data struct {
            Value float64 `json:"value"`
        }
        json.Unmarshal(event.Body, &data)

        state.mu.Lock()
        state.metrics[event.Channel] = data.Value
        state.mu.Unlock()

        fmt.Printf("[Dashboard] %s = %.2f\n", event.Channel, data.Value)
    }),
    kubemq.WithOnError(func(err error) {
        log.Println("[Dashboard] Error:", err)
    }),
)
dashboard_aggregator.py
import json

metrics_state = {}

def on_event(event):
    data = json.loads(event.body.decode("utf-8"))
    metrics_state[event.channel] = data["value"]
    print(f"[Dashboard] {event.channel} = {data['value']}")

client.subscribe_to_events(
    subscription=EventsSubscription(
        channel="metrics.>",
        on_receive_event_callback=on_event,
        on_error_callback=lambda e: print(f"[Dashboard] Error: {e}"),
    ),
    cancel=CancellationToken(),
)
dashboard_aggregator.js
const metricsState = new Map();

client.subscribeToEvents({
  channel: "metrics.>",
  onEvent: (msg) => {
    const data = JSON.parse(Buffer.from(msg.body).toString());
    metricsState.set(msg.channel, data.value);
    console.log(`[Dashboard] ${msg.channel} = ${data.value}`);
  },
  onError: (err) => console.error("[Dashboard] Error:", err.message),
});
DashboardAggregator.java
ConcurrentHashMap<String, Double> metricsState = new ConcurrentHashMap<>();

client.subscribeToEvents(EventsSubscription.builder()
    .channel("metrics.>")
    .onReceiveEventCallback(event -> {
        String body = new String(event.getBody());
        double value = Double.parseDouble(
            body.replaceAll(".*\"value\":(\\d+\\.?\\d*).*", "$1"));
        metricsState.put(event.getChannel(), value);
        System.out.printf("[Dashboard] %s = %.2f%n", event.getChannel(), value);
    })
    .onErrorCallback(err ->
        System.err.println("[Dashboard] Error: " + err.getMessage()))
    .build());
DashboardAggregator.cs
var metricsState = new ConcurrentDictionary<string, double>();

await foreach (var msg in client.SubscribeToEventsAsync(
    new EventsSubscription { Channel = "metrics.>" }))
{
    var data = JsonSerializer.Deserialize<JsonElement>(msg.Body.Span);
    var value = data.GetProperty("value").GetDouble();
    metricsState[msg.Channel] = value;
    Console.WriteLine($"[Dashboard] {msg.Channel} = {value}");
}
DashboardAggregator.kt
val metricsState = ConcurrentHashMap<String, Double>()

client.subscribeToEvents(
    channel = "metrics.>",
    onEvent = { event ->
        val body = String(event.body)
        val value = Regex(""""value":(\d+\.?\d*)""").find(body)?.groupValues?.get(1)?.toDouble() ?: 0.0
        metricsState[event.channel] = value
        println("[Dashboard] ${event.channel} = $value")
    },
    onError = { err -> System.err.println("[Dashboard] Error: ${err.message}") }
)
dashboard_aggregator.cpp
std::map<std::string, double> metricsState;

client.subscribeToEvents("metrics.>", "",
    [&metricsState](const kubemq::Event& event) {
        // Simple JSON value extraction
        auto pos = event.body.find("\"value\":");
        if (pos != std::string::npos) {
            double value = std::stod(event.body.substr(pos + 8));
            metricsState[event.channel] = value;
            std::cout << "[Dashboard] " << event.channel
                      << " = " << value << std::endl;
        }
    },
    [](const std::string& err) {
        std::cerr << "[Dashboard] Error: " << err << std::endl;
    }
);
dashboard_aggregator.rs
use std::collections::HashMap;
use std::sync::{Arc, Mutex};

let metrics_state: Arc<Mutex<HashMap<String, f64>>> = Arc::new(Mutex::new(HashMap::new()));

let state = metrics_state.clone();
let sub = client
    .subscribe_to_events(
        "metrics.>",
        "",
        move |event| {
            let state = state.clone();
            Box::pin(async move {
                let body = String::from_utf8_lossy(&event.body);
                // Extract "value" from the JSON payload.
                if let Some(value) = body
                    .split("\"value\":")
                    .nth(1)
                    .and_then(|s| s.split([',', '}']).next())
                    .and_then(|s| s.trim().parse::<f64>().ok())
                {
                    state.lock().unwrap().insert(event.channel.clone(), value);
                    println!("[Dashboard] {} = {}", event.channel, value);
                }
            })
        },
        None,
    )
    .await?;
dashboard_aggregator.rb
require 'kubemq'
require 'json'

metrics_state = {}
cancel = KubeMQ::CancellationToken.new

sub = KubeMQ::PubSub::EventsSubscription.new(channel: 'metrics.>')
client.subscribe_to_events(
  sub,
  cancellation_token: cancel,
  on_error: ->(e) { warn "[Dashboard] Error: #{e.message}" }
) do |event|
  data = JSON.parse(event.body)
  metrics_state[event.channel] = data['value']
  puts "[Dashboard] #{event.channel} = #{data['value']}"
end
dashboard_aggregator.exs
# Aggregate metrics in an Agent so the wildcard handler can update shared state.
{:ok, metrics_state} = Agent.start_link(fn -> %{} end)

{:ok, sub} =
  KubeMQ.Client.subscribe_to_events(client, "metrics.>",
    on_event: fn event ->
      %{"value" => value} = Jason.decode!(event.body)
      Agent.update(metrics_state, &Map.put(&1, event.channel, value))
      IO.puts("[Dashboard] #{event.channel} = #{value}")
    end,
    on_error: fn err -> IO.warn("[Dashboard] Error: #{err.message}") end
  )

Production Considerations

Was this page helpful?

On this page