Queues
Send, receive, peek, and ack durable FIFO queue messages over CloudEvents HTTP with the KubeMQ CloudEvents connector.
The CloudEvents connector exposes KubeMQ's durable queue pattern over plain HTTP: producers POST a CloudEvent to a queue, and one consumer at a time pulls messages off it with at-least-once delivery.
Overview
A queue is a durable, FIFO work channel. Unlike events (fire-and-forget pub/sub), a queue message is stored until a consumer receives it, and is delivered to exactly one consumer in the group — making queues the right fit for distributing work across a pool of workers.
The connector maps three control operations onto HTTP:
| Operation | Endpoint | Body | Purpose |
|---|---|---|---|
| Send | POST /ce/queue/send | CloudEvent (structured or binary) | Enqueue one message |
| Receive | POST /ce/queue/receive | (ignored) — all params via query string | Pull (and consume) messages |
| Ack all | POST /ce/queue/ack_all | (ignored) — all params via query string | Drain all pending messages atomically |
receive and ack_all are control operations: the request body is ignored and every parameter is supplied as a query string. To inspect messages without consuming them, set is_peek=true on a receive call (peek).
How it works
A producer enqueues CloudEvents; a worker polls the queue and the connector hands each stored message to exactly one consumer.
Messages are durably stored on the queue channel and delivered to one worker at a time.
Send and receive
Send a CloudEvent to a queue with POST /ce/queue/send (returns HTTP 202), then pull it back with POST /ce/queue/receive. The receive call is a control operation — pass channel, client_id, max_messages, and wait_timeout as query parameters.
# Send a message to the queue (structured CloudEvent)
curl -X POST http://localhost:9090/ce/queue/send \
-H "Content-Type: application/cloudevents+json" \
-d '{
"specversion": "1.0",
"type": "com.example.job.submit",
"source": "job-scheduler",
"subject": "work-queue",
"data": {"job_type": "report", "params": {"month": "2026-03"}}
}'
# Receive (consume) up to 10 messages, waiting up to 10s
curl -X POST "http://localhost:9090/ce/queue/receive?channel=work-queue&client_id=worker-1&max_messages=10&wait_timeout=10"using CloudNative.CloudEvents;
using CloudNative.CloudEvents.SystemTextJson;
using System.Net.Http.Headers;
using System.Text.Json;
static string ServerUrl() => Environment.GetEnvironmentVariable("KUBEMQ_CE_URL") ?? "http://localhost:9090";
var base_ = ServerUrl();
var channel = "csharp-ce-queues.basic";
var clientId = "kubemq-ce-csharp-worker";
var formatter = new JsonEventFormatter();
using var httpClient = new HttpClient();
Console.WriteLine($"Sending 3 messages to queue '{channel}':");
for (int i = 1; i <= 3; i++)
{
var ev = new CloudEvent {
Id = Guid.NewGuid().ToString(),
Type = "com.kubemq.examples.queues.task",
Source = new Uri("urn:kubemq-ce-csharp-example"),
Subject = channel,
DataContentType = "application/json",
Data = new { task_id = i, task = "process-item" },
};
var bytes = formatter.EncodeStructuredModeMessage(ev, out var ct);
using var c = new ByteArrayContent(bytes.ToArray());
c.Headers.ContentType = MediaTypeHeaderValue.Parse(ct.ToString());
var r = await httpClient.PostAsync($"{base_}/ce/queue/send", c);
var j = JsonSerializer.Deserialize<JsonElement>(await r.Content.ReadAsStringAsync());
Console.WriteLine($" Sent task {i}: is_error={j.GetProperty("is_error")}");
}
Console.WriteLine("\nReceiving 3 messages:");
for (int i = 0; i < 3; i++)
{
var url = $"{base_}/ce/queue/receive?channel={Uri.EscapeDataString(channel)}&client_id={clientId}&max_messages=1&wait_timeout=5";
var r = await httpClient.PostAsync(url, null);
var j = JsonSerializer.Deserialize<JsonElement>(await r.Content.ReadAsStringAsync());
if (j.GetProperty("is_error").GetBoolean()) {
Console.WriteLine($" Error: {j.GetProperty("message")}"); continue;
}
if (j.TryGetProperty("data", out var data) && data.TryGetProperty("messages", out var msgs))
{
foreach (var m in msgs.EnumerateArray())
Console.WriteLine($" Received: type={m.GetProperty("type")} data={m.GetProperty("data")}");
}
}package main
import (
"encoding/json"
"fmt"
"log"
"net/http"
"os"
"strings"
cloudevents "github.com/cloudevents/sdk-go/v2"
)
func serverURL() string {
if u := os.Getenv("KUBEMQ_CE_URL"); u != "" {
return u
}
return "http://localhost:9090"
}
type CEResponse struct {
IsError bool `json:"is_error"`
Message string `json:"message"`
Data json.RawMessage `json:"data"`
}
type QueueReceiveData struct {
MessagesReceived int `json:"messages_received"`
Messages []map[string]interface{} `json:"messages"`
}
func sendToQueue(base, channel string, i int) {
event := cloudevents.NewEvent()
event.SetType("com.kubemq.examples.queues.task")
event.SetSource("kubemq-ce-go-example")
event.SetSubject(channel)
_ = event.SetData(cloudevents.ApplicationJSON, map[string]interface{}{
"task_id": i, "task": "process-item", "priority": "normal",
})
body, _ := json.Marshal(event)
req, _ := http.NewRequest("POST", base+"/ce/queue/send", strings.NewReader(string(body)))
req.Header.Set("Content-Type", "application/cloudevents+json")
resp, err := http.DefaultClient.Do(req)
if err != nil {
log.Fatal("queue send:", err)
}
defer resp.Body.Close()
var result CEResponse
_ = json.NewDecoder(resp.Body).Decode(&result)
fmt.Printf("Sent task %d: status=%d is_error=%v\n", i, resp.StatusCode, result.IsError)
}
func receiveFromQueue(base, channel, clientID string) {
url := fmt.Sprintf("%s/ce/queue/receive?channel=%s&client_id=%s&max_messages=1&wait_timeout=5",
base, channel, clientID)
req, _ := http.NewRequest("POST", url, nil)
resp, err := http.DefaultClient.Do(req)
if err != nil {
log.Fatal("queue receive:", err)
}
defer resp.Body.Close()
var result CEResponse
_ = json.NewDecoder(resp.Body).Decode(&result)
var data QueueReceiveData
_ = json.Unmarshal(result.Data, &data)
for _, msg := range data.Messages {
fmt.Printf("Received: type=%v data=%v\n", msg["type"], msg["data"])
}
}
func main() {
base := serverURL()
channel := "go-ce-queues.basic"
clientID := "kubemq-ce-go-worker"
for i := 1; i <= 3; i++ {
sendToQueue(base, channel, i)
}
for i := 0; i < 3; i++ {
receiveFromQueue(base, channel, clientID)
}
}import com.fasterxml.jackson.databind.ObjectMapper;
import io.cloudevents.CloudEvent;
import io.cloudevents.core.builder.CloudEventBuilder;
import io.cloudevents.core.format.EventFormat;
import io.cloudevents.core.provider.EventFormatProvider;
import io.cloudevents.jackson.JsonFormat;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.time.OffsetDateTime;
import java.util.List;
import java.util.Map;
import java.util.UUID;
public class Main {
static String serverUrl() {
String u = System.getenv("KUBEMQ_CE_URL");
return (u != null && !u.isEmpty()) ? u : "http://localhost:9090";
}
static final ObjectMapper MAPPER = new ObjectMapper();
public static void main(String[] args) throws Exception {
String base = serverUrl();
String channel = "java-ce-queues.basic";
String clientId = "kubemq-ce-java-worker";
EventFormatProvider.getInstance().registerFormat(new JsonFormat());
EventFormat format = EventFormatProvider.getInstance().resolveFormat(JsonFormat.CONTENT_TYPE);
HttpClient httpClient = HttpClient.newHttpClient();
for (int i = 1; i <= 3; i++) {
CloudEvent event = CloudEventBuilder.v1()
.withId(UUID.randomUUID().toString())
.withType("com.kubemq.examples.queues.task")
.withSource(URI.create("kubemq-ce-java-example"))
.withSubject(channel)
.withDataContentType("application/json")
.withTime(OffsetDateTime.now())
.withData("application/json",
MAPPER.writeValueAsBytes(Map.of("task_id", i, "task", "process-item")))
.build();
httpClient.send(
HttpRequest.newBuilder().uri(URI.create(base + "/ce/queue/send"))
.POST(HttpRequest.BodyPublishers.ofByteArray(format.serialize(event)))
.header("Content-Type", "application/cloudevents+json").build(),
HttpResponse.BodyHandlers.ofString());
}
for (int i = 0; i < 3; i++) {
String url = base + "/ce/queue/receive?channel=" + channel
+ "&client_id=" + clientId + "&max_messages=1&wait_timeout=5";
HttpResponse<String> resp = httpClient.send(
HttpRequest.newBuilder().uri(URI.create(url))
.POST(HttpRequest.BodyPublishers.noBody()).build(),
HttpResponse.BodyHandlers.ofString());
Map<?, ?> r = MAPPER.readValue(resp.body(), Map.class);
Map<?, ?> data = (Map<?, ?>) r.get("data");
List<?> msgs = data != null ? (List<?>) data.get("messages") : List.of();
for (Object m : msgs) {
Map<?, ?> msg = (Map<?, ?>) m;
System.out.println(" Received: type=" + msg.get("type") + " data=" + msg.get("data"));
}
}
}
}import { CloudEvent, HTTP } from 'cloudevents';
function serverUrl() {
return process.env.KUBEMQ_CE_URL ?? 'http://localhost:9090';
}
async function main() {
const base = serverUrl();
const channel = 'js-ce-queues.basic';
const clientId = 'kubemq-ce-js-worker';
const numMessages = 3;
for (let i = 1; i <= numMessages; i++) {
const event = new CloudEvent({
type: 'com.kubemq.examples.queues.task',
source: 'kubemq-ce-js-example',
subject: channel,
datacontenttype: 'application/json',
data: { task_id: i, task: 'process-item' },
});
const msg = HTTP.structured(event);
const resp = await fetch(`${base}/ce/queue/send`, {
method: 'POST',
headers: msg.headers,
body: msg.body,
});
const result = await resp.json();
console.log(` Sent task ${i}: status=${resp.status} is_error=${result.is_error}`);
}
for (let i = 0; i < numMessages; i++) {
const url = `${base}/ce/queue/receive?channel=${encodeURIComponent(channel)}&client_id=${clientId}&max_messages=1&wait_timeout=5`;
const resp = await fetch(url, { method: 'POST' });
const result = await resp.json();
for (const m of result.data?.messages ?? []) {
console.log(` Received: type=${m.type} data=${JSON.stringify(m.data)}`);
}
}
}
main().catch(console.error);import os
import requests
from cloudevents.v1.conversion import to_structured
from cloudevents.v1.http import CloudEvent
def server_url() -> str:
return os.environ.get("KUBEMQ_CE_URL", "http://localhost:9090")
def main() -> None:
base = server_url()
channel = "python-ce-queues.basic"
client_id = "kubemq-ce-python-worker"
num_messages = 3
for i in range(1, num_messages + 1):
event = CloudEvent(
attributes={
"type": "com.kubemq.examples.queues.task",
"source": "kubemq-ce-python-example",
"subject": channel,
"datacontenttype": "application/json",
},
data={"task_id": i, "task": "process-item"},
)
headers, body = to_structured(event)
resp = requests.post(f"{base}/ce/queue/send", data=body,
headers=dict(headers), timeout=10)
result = resp.json()
print(f" Sent task {i}: status={resp.status_code} is_error={result.get('is_error')}")
for _ in range(num_messages):
resp = requests.post(
f"{base}/ce/queue/receive",
params={
"channel": channel,
"client_id": client_id,
"max_messages": 1,
"wait_timeout": 5,
},
timeout=10,
)
result = resp.json()
for msg in result.get("data", {}).get("messages", []):
print(f" Received: type={msg.get('type')} data={msg.get('data')}")
if __name__ == "__main__":
main()require "net/http"
require "uri"
require "json"
require "securerandom"
require "cloud_events"
def server_url = ENV.fetch("KUBEMQ_CE_URL", "http://localhost:9090")
base = server_url
channel = "ruby-ce-queues.basic"
client_id = "kubemq-ce-ruby-worker"
sdk = CloudEvents::HttpBinding.default
(1..3).each do |i|
event = CloudEvents::Event::V1.new(
id: SecureRandom.uuid, type: "com.kubemq.examples.queues.task",
source: URI("urn:kubemq-ce-ruby-example"), subject: channel,
spec_version: "1.0",
data_content_type: CloudEvents::ContentType.new("application/json"),
data: JSON.generate({ task_id: i, task: "process-item" })
)
enc_headers, enc_body = sdk.encode_event(event, structured_format: "json")
uri = URI("#{base}/ce/queue/send")
Net::HTTP.start(uri.host, uri.port) do |http|
req = Net::HTTP::Post.new(uri)
enc_headers.each { |k, v| req[k] = v }
req.body = enc_body
res = http.request(req)
r = JSON.parse(res.body)
puts " Sent task #{i}: status=#{res.code} is_error=#{r['is_error']}"
end
end
3.times do
uri = URI("#{base}/ce/queue/receive")
uri.query = URI.encode_www_form(channel: channel, client_id: client_id,
max_messages: 1, wait_timeout: 5)
Net::HTTP.start(uri.host, uri.port) do |http|
res = http.request(Net::HTTP::Post.new(uri))
r = JSON.parse(res.body)
msgs = r.dig('data', 'messages') || []
msgs.each { |m| puts " Received: type=#{m['type']} data=#{m['data']}" }
end
enduse cloudevents::{EventBuilder, EventBuilderV10};
use reqwest::Client;
use serde_json::{json, Value};
use std::env;
use uuid::Uuid;
fn server_url() -> String {
env::var("KUBEMQ_CE_URL").unwrap_or_else(|_| "http://localhost:9090".to_string())
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let base = server_url();
let client = Client::new();
let channel = "rust-ce-queues.basic";
let client_id = "kubemq-ce-rust-worker";
for i in 1..=3u32 {
let event = EventBuilderV10::new()
.id(Uuid::new_v4().to_string())
.ty("com.kubemq.examples.queues.task")
.source("urn:kubemq-ce-rust-example")
.subject(channel)
.data("application/json", json!({"task_id": i, "task": "process-item"}))
.build()?;
let body = serde_json::to_string(&event)?;
let resp = client.post(format!("{}/ce/queue/send", base))
.header("Content-Type", "application/cloudevents+json")
.body(body).send().await?;
let r: Value = resp.json().await?;
println!(" Sent task {}: is_error={}", i, r["is_error"]);
}
for _ in 0..3 {
let url = format!(
"{}/ce/queue/receive?channel={}&client_id={}&max_messages=1&wait_timeout=5",
base, channel, client_id
);
let resp = client.post(&url).send().await?;
let r: Value = resp.json().await?;
let msgs = r["data"]["messages"].as_array().cloned().unwrap_or_default();
for m in &msgs {
println!(" Received: type={} data={}", m["type"], m["data"]);
}
}
Ok(())
}Receive parameters
| Parameter | Type | Default | Description |
|---|---|---|---|
channel | string | (required) | Queue channel name |
client_id | string | (required) | Client identifier (overridden by auth claims when auth is enabled) |
max_messages | int | 1 | Maximum number of messages to receive (1–1000) |
wait_timeout | int | 5 | Wait timeout in seconds |
is_peek | bool | false | If true, peek at messages without consuming them |
A successful receive returns HTTP 200. Each message that carries CE tags is reconstructed as a CloudEvent JSON object; messages without CE tags are returned in plain KubeMQ format:
{
"is_error": false,
"message": "OK",
"data": {
"messages_received": 1,
"messages": [
{
"specversion": "1.0",
"type": "com.example.job.submit",
"source": "job-scheduler",
"subject": "work-queue",
"id": "550e8400-e29b-41d4-a716-446655440000",
"time": "2026-03-29T10:30:00Z",
"data": {"job_type": "report", "params": {"month": "2026-03"}}
}
]
}
}Peek
Peeking inspects queued messages without consuming them — pass is_peek=true on a receive call. The same messages remain available for a later consuming receive, which makes peek useful for monitoring queue depth or previewing work before committing to it.
# Peek: inspect up to 10 messages without removing them
curl -X POST "http://localhost:9090/ce/queue/receive?channel=work-queue&client_id=worker-1&max_messages=10&wait_timeout=3&is_peek=true"
# A normal receive (is_peek omitted) consumes the same messages
curl -X POST "http://localhost:9090/ce/queue/receive?channel=work-queue&client_id=worker-1&max_messages=10&wait_timeout=3"func receiveOrPeek(base, channel, clientID string, isPeek bool, label string) int {
isPeekStr := "false"
if isPeek {
isPeekStr = "true"
}
url := fmt.Sprintf(
"%s/ce/queue/receive?channel=%s&client_id=%s&max_messages=10&wait_timeout=3&is_peek=%s",
base, channel, clientID, isPeekStr)
req, _ := http.NewRequest("POST", url, nil)
resp, err := http.DefaultClient.Do(req)
if err != nil {
log.Fatal("receive:", err)
}
defer resp.Body.Close()
var result CEResponse
_ = json.NewDecoder(resp.Body).Decode(&result)
var data QueueReceiveData
_ = json.Unmarshal(result.Data, &data)
fmt.Printf("[%s] messages_received=%d\n", label, data.MessagesReceived)
return data.MessagesReceived
}
// Peek twice (messages stay), then consume.
receiveOrPeek(base, channel, "go-peek-client", true, "peek #1")
receiveOrPeek(base, channel, "go-peek-client", true, "peek #2")
receiveOrPeek(base, channel, "go-peek-client", false, "consume")async function receiveOrPeek(base, channel, clientId, isPeek, label) {
const url = `${base}/ce/queue/receive?channel=${encodeURIComponent(channel)}&client_id=${clientId}&max_messages=10&wait_timeout=3&is_peek=${isPeek}`;
const resp = await fetch(url, { method: 'POST' });
const result = await resp.json();
const n = result.data?.messages_received ?? 0;
console.log(`[${label}] messages_received=${n}`);
return n;
}
// Peek twice (messages stay), then consume.
await receiveOrPeek(base, channel, 'js-peek-client', true, 'peek #1');
await receiveOrPeek(base, channel, 'js-peek-client', true, 'peek #2');
await receiveOrPeek(base, channel, 'js-peek-client', false, 'consume');def receive_or_peek(base, channel, client_id, is_peek, label):
resp = requests.post(
f"{base}/ce/queue/receive",
params={
"channel": channel,
"client_id": client_id,
"max_messages": 10,
"wait_timeout": 3,
"is_peek": "true" if is_peek else "false",
},
timeout=10,
)
result = resp.json()
n = result.get("data", {}).get("messages_received", 0)
print(f"[{label}] messages_received={n}")
return n
# Peek twice (messages stay), then consume.
receive_or_peek(base, channel, "python-peek-client", True, "peek #1")
receive_or_peek(base, channel, "python-peek-client", True, "peek #2")
receive_or_peek(base, channel, "python-peek-client", False, "consume")The query parameter is spelled is_peek on the wire — setting is_peek=true performs a non-destructive peek (messages stay queued), while omitting it (or is_peek=false) performs a normal consuming receive. The internal Go struct field is named IsPeak, but that name never appears on the wire.
Ack all
POST /ce/queue/ack_all acknowledges (drains) all pending messages in a queue atomically, without receiving them one by one. This is a control operation — the body is ignored and channel, client_id, and wait_timeout come from query parameters. Use it to clear a backlog or reset a work channel.
curl -X POST "http://localhost:9090/ce/queue/ack_all?channel=work-queue&client_id=worker-1&wait_timeout=10"// Drain the queue atomically.
ackURL := fmt.Sprintf("%s/ce/queue/ack_all?channel=%s&client_id=%s&wait_timeout=5",
base, channel, clientID)
req, _ := http.NewRequest("POST", ackURL, nil)
resp, err := http.DefaultClient.Do(req)
if err != nil {
log.Fatal("ack_all:", err)
}
defer resp.Body.Close()
var ackResult CEResponse
_ = json.NewDecoder(resp.Body).Decode(&ackResult)
fmt.Printf("ack_all: is_error=%v message=%s\n", ackResult.IsError, ackResult.Message)# Drain the queue atomically.
ack_resp = requests.post(
f"{base}/ce/queue/ack_all",
params={"channel": channel, "client_id": client_id, "wait_timeout": 5},
timeout=10,
)
result = ack_resp.json()
print(f"ack_all: is_error={result.get('is_error')} message={result.get('message')}")Ack-all parameters
| Parameter | Type | Default | Description |
|---|---|---|---|
channel | string | (required) | Queue channel name |
client_id | string | (required) | Client identifier (overridden by auth claims when auth is enabled) |
wait_timeout | int | 5 | Wait timeout in seconds |
A successful ack returns HTTP 200 with the result in data.
send, receive, and ack_all are synchronous POST endpoints and are subject to the connector's TimeoutSeconds (default 60s). A request that exceeds the timeout returns HTTP 504. See Configuration.
Related
Was this page helpful?