Stream API (Upstream/Downstream)
Use bidirectional streaming for continuous queue message sending and receiving.
What You Will Build
A streaming producer that sends messages continuously through an upstream connection, and a streaming consumer that receives messages through a downstream connection with range-based acknowledgment.
One long-lived upstream stream feeds the queue; a downstream stream delivers batches that the consumer settles with a dotted range acknowledgment.
Prerequisites
- KubeMQ server running on
localhost:50000 - SDK installed (Getting Started)
Upstream: Stream Send
Open a persistent bidirectional stream for high-throughput sending. Each message gets an individual response confirming receipt.
stream, err := client.UpstreamQueue(ctx)
if err != nil {
log.Fatal(err)
}
for i := 0; i < 100; i++ {
body := fmt.Sprintf(`{"orderId":"ORD-%04d","item":"widget"}`, i)
result, err := stream.Send(ctx, kubemq.NewQueueMessage().
SetChannel("orders.stream").
SetBody([]byte(body)),
)
if err != nil {
log.Printf("Send error: %v", err)
continue
}
fmt.Printf("Streamed: id=%s\n", result.MessageID)
}
stream.Close()stream = client.upstream_queue()
for i in range(100):
result = stream.send(
QueueMessage(
channel="orders.stream",
body=f'{{"orderId":"ORD-{i:04d}","item":"widget"}}'.encode(),
)
)
print(f"Streamed: id={result.id}")
stream.close()const upstream = client.createQueueUpstream();
for (let i = 0; i < 100; i++) {
const result = await upstream.send([
createQueueMessage({
channel: 'orders.stream',
body: JSON.stringify({ orderId: `ORD-${String(i).padStart(4, '0')}`, item: 'widget' }),
}),
]);
console.log(`Streamed: id=${result.results[0].messageId}`);
}
upstream.close();QueueUpstream stream = client.upstreamQueue();
for (int i = 0; i < 100; i++) {
SendQueueMessageResult result = stream.send(QueueMessage.builder()
.channel("orders.stream")
.body(String.format("{\"orderId\":\"ORD-%04d\",\"item\":\"widget\"}", i).getBytes())
.build());
System.out.printf("Streamed: id=%s%n", result.getMessageId());
}
stream.close();var stream = await client.UpstreamQueueAsync();
for (int i = 0; i < 100; i++)
{
var result = await stream.SendAsync(new QueueMessage
{
Channel = "orders.stream",
Body = Encoding.UTF8.GetBytes($"{{\"orderId\":\"ORD-{i:D4}\",\"item\":\"widget\"}}")
});
Console.WriteLine($"Streamed: id={result.MessageId}");
}
await stream.CloseAsync();val stream = client.upstreamQueue()
for (i in 0 until 100) {
val result = stream.send(QueueMessage(
channel = "orders.stream",
body = """{"orderId":"ORD-${"%04d".format(i)}","item":"widget"}""".toByteArray()
))
println("Streamed: id=${result.messageId}")
}
stream.close()auto stream = client.upstreamQueue();
for (int i = 0; i < 100; i++) {
kubemq::QueueMessage msg;
msg.channel = "orders.stream";
msg.body = "{\"orderId\":\"ORD-" + std::to_string(i) + "\",\"item\":\"widget\"}";
auto result = stream.send(msg);
std::cout << "Streamed: id=" << result.messageId << std::endl;
}
stream.close();let mut upstream = client.queue_upstream().await?;
// Send a batch of messages over the persistent upstream stream.
let messages: Vec<QueueMessage> = (0..100)
.map(|i| {
QueueMessageBuilder::new()
.channel("orders.stream")
.body(format!(r#"{{"orderId":"ORD-{:04}","item":"widget"}}"#, i).into_bytes())
.build()
})
.collect();
upstream.send("orders-batch-001", messages).await?;
if let Some(result) = upstream.results().recv().await {
println!(
"Streamed: ref_id={}, is_error={}, items={}",
result.ref_request_id,
result.is_error,
result.results.len()
);
}
upstream.close();sender = client.create_upstream_sender
100.times do |i|
msg = KubeMQ::Queues::QueueMessage.new(
channel: "orders.stream",
body: %({"orderId":"ORD-#{format('%04d', i)}","item":"widget"})
)
results = sender.publish(msg)
results.each { |r| puts "Streamed: id=#{r.id}, error?=#{r.error?}" }
end
sender.close{:ok, handle} = KubeMQ.Client.queue_upstream(client)
# Send a batch of messages over the persistent upstream stream.
messages =
for i <- 0..99 do
body = ~s({"orderId":"ORD-#{String.pad_leading(Integer.to_string(i), 4, "0")}","item":"widget"})
KubeMQ.QueueMessage.new(channel: "orders.stream", body: body)
end
case KubeMQ.QueueUpstreamHandle.send(handle, messages) do
{:ok, results} ->
Enum.each(results, fn r ->
IO.puts("Streamed: id=#{r.message_id}, error: #{r.is_error}")
end)
{:error, err} ->
IO.puts("Stream send failed: #{err.message}")
end
KubeMQ.QueueUpstreamHandle.close(handle)Downstream: Stream Receive
Open a persistent downstream stream for continuous message consumption.
stream, err := client.DownstreamQueue(ctx, &kubemq.DownstreamRequest{
Channel: "orders.stream",
MaxItems: 10,
WaitTimeoutSeconds: 5,
AutoAck: false,
})
if err != nil {
log.Fatal(err)
}
for resp := range stream.ResponseCh {
for _, m := range resp.Messages {
fmt.Printf("Received: %s\n", string(m.Message.Body))
}
resp.AckAll()
}def on_message(response):
for msg in response.messages:
print(f"Received: {msg.body.decode('utf-8')}")
msg.ack()
def on_error(err):
print(f"Stream error: {err}")
stream = client.downstream_queue(
channel="orders.stream",
max_messages=10,
wait_timeout_in_seconds=5,
on_message_callback=on_message,
on_error_callback=on_error,
)const stream = client.streamQueueMessages({
channel: 'orders.stream',
maxMessages: 10,
waitTimeoutSeconds: 5,
autoAck: false,
});
stream.onMessages((messages) => {
for (const msg of messages) {
console.log('Received:', new TextDecoder().decode(msg.body));
}
// Acknowledge the whole batch at once.
stream.ackAll();
});
stream.onError((err) => console.error('Stream error:', err.message));client.downstreamQueue(DownstreamRequest.builder()
.channel("orders.stream")
.maxMessages(10)
.waitTimeoutSeconds(5)
.onMessage(response -> {
for (QueueMessageReceived msg : response.getMessages()) {
System.out.println("Received: " + new String(msg.getBody()));
msg.ack();
}
})
.onError(err -> System.err.println("Stream error: " + err.getMessage()))
.build());await foreach (var response in client.DownstreamQueueAsync(new DownstreamRequest
{
Channel = "orders.stream",
MaxMessages = 10,
WaitTimeoutSeconds = 5,
}))
{
foreach (var msg in response.Messages)
{
Console.WriteLine($"Received: {Encoding.UTF8.GetString(msg.Body.Span)}");
await msg.AckAsync();
}
}client.downstreamQueue(
channel = "orders.stream",
maxMessages = 10,
waitTimeoutSeconds = 5,
onMessage = { response ->
for (msg in response.messages) {
println("Received: ${String(msg.body)}")
msg.ack()
}
},
onError = { err -> System.err.println("Stream error: ${err.message}") }
)client.downstreamQueue("orders.stream", 10, 5,
[](const auto& response) {
for (const auto& msg : response.messages) {
std::cout << "Received: " << msg.body << std::endl;
msg.ack();
}
},
[](const std::string& err) {
std::cerr << "Stream error: " << err << std::endl;
}
);let mut receiver = client.new_queue_downstream_receiver().await?;
// Poll a batch over the persistent downstream stream.
let poll = PollRequest {
channel: "orders.stream".to_string(),
max_items: 10,
wait_timeout_seconds: 5,
auto_ack: false,
};
let response = receiver.poll(poll).await?;
for msg in &response.messages {
println!("Received: {}", String::from_utf8_lossy(&msg.body));
}
// Acknowledge the whole batch (range ack) in one call.
response.ack_all().await?;
receiver.close().await?;receiver = client.create_downstream_receiver
request = KubeMQ::Queues::QueuePollRequest.new(
channel: "orders.stream",
max_items: 10,
wait_timeout: 5
)
response = receiver.poll(request)
if response.error?
puts "Stream error: #{response.error}"
else
response.messages.each do |msg|
puts "Received: #{msg.body}"
msg.ack
end
end
receiver.close# Poll a batch over the downstream stream.
case KubeMQ.Client.poll_queue(client,
channel: "orders.stream",
max_items: 10,
wait_timeout: 5_000
) do
{:ok, poll} ->
Enum.each(poll.messages, fn msg ->
IO.puts("Received: #{msg.body}")
end)
# Acknowledge the whole batch (range ack) in one call.
{:ok, _} = KubeMQ.PollResponse.ack_all(poll)
{:error, err} ->
IO.puts("Stream error: #{err.message}")
endStream vs Polling Comparison
| Feature | Stream API | Polling (PollQueue) |
|---|---|---|
| Connection | Persistent bidirectional | Request/response per poll |
| Latency | Lower (always connected) | Higher (new request each time) |
| Throughput | Higher | Lower |
| Resource usage | Holds connection open | Releases between polls |
| Best for | High-volume continuous processing | Periodic batch processing |
Next Steps
Was this page helpful?