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.
| Operation | AMQP action | KubeMQ mapping |
|---|---|---|
| Produce | Attach a sender to events/<ch>, TRANSFER (pre-settle for fire-and-forget) | SendEvents (Store=false) |
| Consume | Attach a receiver to events/<ch>, grant credit | Pre-settled fan-out delivery (at-most-once) |
| Consumer group | Set link property x-opt-kubemq-group on the receiver | Load-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.
Related
Was this page helpful?