KubeMQ
ConnectorsAMQP 1.0Concepts

Events

Fire-and-forget pub/sub over AMQP 1.0 — pre-settled at-most-once fan-out, standing credit, and x-opt-kubemq-group consumer groups on the KubeMQ Events pattern.

Events are fire-and-forget pub/sub over the AMQP 1.0 connector. Attach a sender or receiver to a node whose address begins with events/<channel> and the connector binds the link to the KubeMQ Events pattern. Every active subscriber receives a copy of every message; there is no persistence and no replay.

Overview

The events/ prefix selects the Events pattern (longest-prefix match, evaluated before events-store/). A producer attaches a sender to events/<ch>; each TRANSFER becomes a KubeMQ SendEvents with Store=false. A consumer attaches a receiver to the same node and grants credit; deliveries are always pre-settled (settled=true) — at-most-once, with no DISPOSITION round-trip.

Use Events for real-time notifications, telemetry, and broadcast where a missed message is acceptable. When you need durability and replay, use Events Store instead.

OperationAMQP actionKubeMQ mapping
ProduceAttach a sender to events/<ch>, TRANSFER (pre-settle for fire-and-forget)SendEvents (Store=false)
ConsumeAttach a receiver to events/<ch>, grant creditPre-settled fan-out delivery (at-most-once)
Consumer groupSet link property x-opt-kubemq-group on the receiverLoad-balanced subset of the stream

How it works

A published event fans out to every connected receiver on the channel. Receivers that share an x-opt-kubemq-group link property split the stream load-balanced; a receiver with no group is a plain fan-out subscriber.

Each event is copied to every plain subscriber; members sharing an x-opt-kubemq-group split the stream between them.

Events at 0 credit are silently dropped. A message that arrives at a receiver whose link credit is 0 is discarded with no error and no DISPOSITION — that is what at-most-once means here. Grant a standing credit and replenish it eagerly, and subscribe before you publish (there is no replay to catch up from). The connector counts every drop in kubemq_amqp10_events_dropped_no_credit_total.

Publish and subscribe

Each example below subscribes first (a receiver with a large standing credit), waits ~750 ms for the connector's subscription pump to go live, then publishes pre-settled events and drains them. Every client reads the broker endpoint from KUBEMQ_AMQP_URL (default amqp://localhost:5672).

package main

import (
	"context"
	"fmt"
	"log"
	"os"
	"time"

	amqp "github.com/Azure/go-amqp"
)

const channel = "amqp10.examples.pubsub"
const total = 20
const standingCredit = 100 // never let credit reach 0 — a 0-credit event is dropped

func amqpURL() string {
	if v := os.Getenv("KUBEMQ_AMQP_URL"); v != "" {
		return v
	}
	return "amqp://localhost:5672"
}

func main() {
	ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
	defer cancel()

	addr := "events/" + channel // events/ prefix → KubeMQ Events pattern

	conn, err := amqp.Dial(ctx, amqpURL(), nil)
	if err != nil {
		log.Fatalf("dial: %v", err)
	}
	defer func() { _ = conn.Close() }()
	session, err := conn.NewSession(ctx, nil)
	if err != nil {
		log.Fatalf("new session: %v", err)
	}

	// 1. SUBSCRIBE FIRST with standing credit. Events have no replay — a publish
	//    that beats the subscription is lost forever.
	receiver, err := session.NewReceiver(ctx, addr, &amqp.ReceiverOptions{Credit: standingCredit})
	if err != nil {
		log.Fatalf("new receiver: %v", err)
	}
	// The attach reply confirms the link, not that the subscription pump is live.
	time.Sleep(750 * time.Millisecond)

	// 2. PUBLISH pre-settled (fire-and-forget) — no DISPOSITION to await.
	sender, err := session.NewSender(ctx, addr, &amqp.SenderOptions{
		SettlementMode: amqp.SenderSettleModeSettled.Ptr(),
	})
	if err != nil {
		log.Fatalf("new sender: %v", err)
	}
	for i := 0; i < total; i++ {
		if err := sender.Send(ctx, amqp.NewMessage([]byte(fmt.Sprintf("event-%03d", i))), nil); err != nil {
			log.Fatalf("publish: %v", err)
		}
	}
	_ = sender.Close(ctx)

	// 3. RECEIVE. Standing credit drains every event; accept is a no-op on
	//    pre-settled fan-out but harmless.
	seen := make(map[string]struct{}, total)
	for len(seen) < total {
		msg, err := receiver.Receive(ctx, nil)
		if err != nil {
			log.Fatalf("receive: %v", err)
		}
		_ = receiver.AcceptMessage(ctx, msg)
		seen[string(msg.GetData())] = struct{}{}
	}
	fmt.Printf("received all %d events\n", len(seen))
	_ = receiver.Close(ctx)
}
import os
import time

from proton import Message
from proton.reactor import AtMostOnce
from proton.utils import BlockingConnection

CHANNEL = "amqp10.examples.pubsub"
TOTAL = 20
STANDING_CREDIT = 100  # never let credit reach 0 — a 0-credit event is dropped


def amqp_url() -> str:
    return os.environ.get("KUBEMQ_AMQP_URL", "amqp://localhost:5672")


def accept_if_unsettled(receiver) -> None:
    # Events fan-out deliveries are pre-settled, so accept() on a settled delivery
    # raises IndexError. This makes accept a true no-op on pre-settled pub/sub.
    if receiver.fetcher.unsettled:
        receiver.accept()


def main() -> None:
    addr = "events/" + CHANNEL  # events/ prefix → KubeMQ Events pattern
    conn = BlockingConnection(amqp_url())
    try:
        # 1. SUBSCRIBE FIRST with standing credit (events have no replay).
        receiver = conn.create_receiver(addr, credit=STANDING_CREDIT)
        time.sleep(0.75)  # let the subscription pump go live before publishing

        # 2. PUBLISH pre-settled (AtMostOnce) — fire-and-forget, no DISPOSITION.
        sender = conn.create_sender(addr, options=AtMostOnce())
        for i in range(TOTAL):
            sender.send(Message(body=f"event-{i:03d}"))
        sender.close()

        # 3. RECEIVE. Standing credit drains every event.
        seen: set[str] = set()
        while len(seen) < TOTAL:
            msg = receiver.receive(timeout=30.0)
            accept_if_unsettled(receiver)
            seen.add(str(msg.body))
        print(f"received all {len(seen)} events")
        receiver.close()
    finally:
        conn.close()


if __name__ == "__main__":
    main()
import java.util.HashSet;
import java.util.Set;

import javax.jms.Connection;
import javax.jms.DeliveryMode;
import javax.jms.Message;
import javax.jms.MessageConsumer;
import javax.jms.MessageProducer;
import javax.jms.Session;
import javax.jms.Topic;

import org.apache.qpid.jms.JmsConnectionFactory;

public final class Main {
    private static final String CHANNEL = "amqp10.examples.pubsub";
    private static final int TOTAL = 20;

    public static void main(String[] args) throws Exception {
        String url = System.getenv().getOrDefault("KUBEMQ_AMQP_URL", "amqp://localhost:5672");
        String address = "events/" + CHANNEL; // events/ prefix → KubeMQ Events pattern

        JmsConnectionFactory factory = new JmsConnectionFactory(url);
        try (Connection connection = factory.createConnection()) {
            connection.start();
            try (Session session = connection.createSession(false, Session.AUTO_ACKNOWLEDGE)) {
                Topic topic = session.createTopic(address);

                // 1. SUBSCRIBE FIRST (Qpid JMS grants a standing prefetch credit).
                try (MessageConsumer consumer = session.createConsumer(topic)) {
                    Thread.sleep(750); // let the subscription pump go live

                    // 2. PUBLISH NON_PERSISTENT (fire-and-forget, pre-settled).
                    try (MessageProducer producer = session.createProducer(topic)) {
                        producer.setDeliveryMode(DeliveryMode.NON_PERSISTENT);
                        for (int i = 0; i < TOTAL; i++) {
                            producer.send(session.createTextMessage(String.format("event-%03d", i)));
                        }
                    }

                    // 3. RECEIVE. The connector re-emits the body as a Data section;
                    //    getBody(String.class) decodes either type as UTF-8.
                    Set<String> seen = new HashSet<>();
                    while (seen.size() < TOTAL) {
                        Message msg = consumer.receive(30_000);
                        if (msg == null) throw new IllegalStateException("timed out");
                        seen.add(msg.getBody(String.class));
                    }
                    System.out.printf("received all %d events%n", seen.size());
                }
            }
        }
    }
}
using System.Text;
using Amqp;
using Amqp.Framing;

const string channel = "amqp10.examples.pubsub";
const int total = 20;
const int standingCredit = 100; // never let credit reach 0 — a 0-credit event is dropped

static string AmqpUrl() =>
    Environment.GetEnvironmentVariable("KUBEMQ_AMQP_URL") is { Length: > 0 } v
        ? v
        : "amqp://localhost:5672";

var addr = "events/" + channel; // events/ prefix → KubeMQ Events pattern

var connection = await Connection.Factory.CreateAsync(new Address(AmqpUrl()));
try
{
    var session = new Session(connection);

    // 1. SUBSCRIBE FIRST with standing credit (autoRestore replenishes on settle).
    var receiver = new ReceiverLink(session, "pubsub-receiver", addr);
    receiver.SetCredit(standingCredit, autoRestore: true);
    await Task.Delay(750); // let the subscription pump go live before publishing

    // 2. PUBLISH pre-settled — SndSettleMode.Settled marks every TRANSFER settled.
    var senderAttach = new Attach
    {
        Source = new Source(),
        Target = new Target { Address = addr },
        SndSettleMode = SenderSettleMode.Settled,
    };
    var sender = new SenderLink(session, "pubsub-sender", senderAttach, null);
    for (var i = 0; i < total; i++)
    {
        var message = new Message { BodySection = new Data { Binary = Encoding.UTF8.GetBytes($"event-{i:D3}") } };
        sender.Send(message, TimeSpan.FromSeconds(15));
    }

    // 3. RECEIVE. Standing credit drains every event; Accept is a no-op here.
    var seen = new HashSet<string>();
    while (seen.Count < total)
    {
        var message = receiver.Receive(TimeSpan.FromSeconds(30))
            ?? throw new InvalidOperationException("receive timed out");
        receiver.Accept(message);
        seen.Add(BodyString(message));
    }
    Console.WriteLine($"received all {seen.Count} events");

    await sender.CloseAsync();
    await receiver.CloseAsync();
    await session.CloseAsync();
}
finally
{
    await connection.CloseAsync();
}

static string BodyString(Message message) => message.BodySection switch
{
    Data d => Encoding.UTF8.GetString(d.Binary),
    AmqpValue { Value: byte[] bytes } => Encoding.UTF8.GetString(bytes),
    AmqpValue { Value: string str } => str,
    AmqpValue v => v.Value?.ToString() ?? string.Empty,
    _ => string.Empty,
};
import {
  Connection,
  ReceiverEvents,
  type EventContext,
  type Receiver,
} from "rhea-promise";

const channel = "amqp10.examples.pubsub";
const total = 20;
const standingCredit = 100; // never let credit reach 0 — a 0-credit event is dropped

function sleep(ms: number): Promise<void> {
  return new Promise((resolve) => setTimeout(resolve, ms));
}

function bodyToString(body: unknown): string {
  return Buffer.isBuffer(body) ? body.toString("utf8") : String(body);
}

async function main(): Promise<void> {
  const url = new URL(process.env["KUBEMQ_AMQP_URL"] ?? "amqp://localhost:5672");
  const address = `events/${channel}`; // events/ prefix → KubeMQ Events pattern

  const connection = new Connection({
    host: url.hostname,
    port: url.port ? Number(url.port) : 5672,
    container_id: `kubemq-amqp10-js-pubsub-${process.pid}`,
    reconnect: false,
  });
  await connection.open();

  try {
    // 1. SUBSCRIBE FIRST. Register the handler before granting credit so no
    //    early delivery is missed (events have no replay).
    const receiver = await connection.createReceiver({
      source: { address },
      credit_window: 0,
      autoaccept: false,
      autosettle: false,
    });

    const seen = new Set<string>();
    const received = drainEvents(receiver, (ctx) => {
      ctx.delivery?.accept(); // no-op for pre-settled fan-out, but harmless
      seen.add(bodyToString(ctx.message?.body));
      return seen.size >= total;
    }, standingCredit, 30_000);

    await sleep(750); // let the subscription pump go live before publishing

    // 2. PUBLISH pre-settled (snd_settle_mode: 1) — fire-and-forget.
    const sender = await connection.createSender({
      target: { address },
      snd_settle_mode: 1,
      autosettle: true,
    });
    for (let i = 0; i < total; i++) {
      sender.send({ body: `event-${String(i).padStart(3, "0")}` });
    }
    await sender.close();

    await received;
    console.log(`received all ${seen.size} events`);
    await receiver.close();
  } finally {
    await connection.close();
  }
}

// Grants standing credit and tops it back up as messages arrive so the
// subscriber is never starved (a 0-credit event is silently dropped).
function drainEvents(
  receiver: Receiver,
  onMessage: (ctx: EventContext) => boolean,
  credit: number,
  timeoutMs: number,
): Promise<void> {
  return new Promise<void>((resolve, reject) => {
    const timer = setTimeout(() => {
      receiver.removeListener(ReceiverEvents.message, handler);
      reject(new Error("timed out waiting for events"));
    }, timeoutMs);
    const handler = (ctx: EventContext): void => {
      if (onMessage(ctx)) {
        clearTimeout(timer);
        receiver.removeListener(ReceiverEvents.message, handler);
        resolve();
        return;
      }
      receiver.addCredit(1); // replenish so standing credit never drains to 0
    };
    receiver.on(ReceiverEvents.message, handler);
    receiver.addCredit(credit);
  });
}

main().catch((err) => {
  console.error(err);
  process.exit(1);
});
use std::collections::HashSet;
use std::time::Duration;

use fe2o3_amqp::link::delivery::Delivery;
use fe2o3_amqp::link::receiver::CreditMode;
use fe2o3_amqp::{Connection, Receiver, Sender, Session};
use fe2o3_amqp_types::definitions::SenderSettleMode;
use fe2o3_amqp_types::messaging::{Body, Message};
use fe2o3_amqp_types::primitives::Value;

const CHANNEL: &str = "amqp10.examples.pubsub";
const TOTAL: usize = 20;
const STANDING_CREDIT: u32 = 100; // never let credit reach 0 — a 0-credit event is dropped

fn amqp_url() -> String {
    std::env::var("KUBEMQ_AMQP_URL").unwrap_or_else(|_| "amqp://localhost:5672".to_string())
}

fn body_string(msg: &Message<Body<Value>>) -> String {
    let bytes = match &msg.body {
        Body::Data(batch) => batch.iter().flat_map(|d| d.0.to_vec()).collect(),
        Body::Value(v) => match &v.0 {
            Value::Binary(b) => b.to_vec(),
            Value::String(s) => s.clone().into_bytes(),
            other => format!("{other:?}").into_bytes(),
        },
        _ => Vec::new(),
    };
    String::from_utf8_lossy(&bytes).into_owned()
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let addr = format!("events/{CHANNEL}"); // events/ prefix → KubeMQ Events pattern

    let mut connection = Connection::open("amqp10-examples-pubsub", amqp_url().as_str()).await?;
    let mut session = Session::begin(&mut connection).await?;

    // 1. SUBSCRIBE FIRST with standing credit (CreditMode::Auto auto-replenishes).
    let mut receiver = Receiver::builder()
        .name("basic-pubsub-receiver")
        .source(addr.as_str())
        .credit_mode(CreditMode::Auto(STANDING_CREDIT))
        .attach(&mut session)
        .await?;
    tokio::time::sleep(Duration::from_millis(750)).await; // let the pump go live

    // 2. PUBLISH pre-settled (SenderSettleMode::Settled) — fire-and-forget.
    let mut sender = Sender::builder()
        .name("basic-pubsub-sender")
        .target(addr.as_str())
        .sender_settle_mode(SenderSettleMode::Settled)
        .attach(&mut session)
        .await?;
    for i in 0..TOTAL {
        sender.send(format!("event-{i:03}")).await?;
    }
    sender.close().await?;

    // 3. RECEIVE. Standing credit drains every event.
    let mut seen: HashSet<String> = HashSet::with_capacity(TOTAL);
    while seen.len() < TOTAL {
        let delivery: Delivery<Body<Value>> = receiver.recv().await?;
        let _ = receiver.accept(&delivery).await; // no-op on pre-settled, harmless
        seen.insert(body_string(delivery.message()));
    }
    println!("received all {} events", seen.len());

    receiver.close().await?;
    session.end().await?;
    connection.close().await?;
    Ok(())
}

Consumer groups

A receiver with no group is a plain fan-out subscriber — it receives every event (the KubeMQ default). Set the link property x-opt-kubemq-group on a receiver's ATTACH to join a consumer group: within one group the stream is load-balanced across members (each message goes to exactly one member), while different groups each get the full stream independently.

events/orders
   ├── group "g1": receiver-A ┐  (split — no duplicate within g1)
   │               receiver-B ┘
   └── group "g2": receiver-C    (full stream)

In Go, set it on the receiver's link properties at attach:

receiver, err := session.NewReceiver(ctx, "events/orders", &amqp.ReceiverOptions{
	Credit:     100,
	Properties: map[string]any{"x-opt-kubemq-group": "g1"},
})

The x-opt-kubemq-group property is honored on events, events-store, and the RPC consume patterns. The Go, Python, C#, JavaScript, and Rust clients can all set it on the receiver's ATTACH.

Qpid JMS (Java) cannot join a consumer group today. The connector advertises no SHARED-SUBS capability, so createSharedConsumer / createSharedDurableConsumer throws "Remote peer does not support shared subscriptions", and Qpid JMS exposes no API to set the x-opt-kubemq-group link property directly. Java is fan-out only on Events; the other five languages support groups fully.

Was this page helpful?

On this page