KubeMQ
LearnQueuesTutorials

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

Upstream: Stream Send

Open a persistent bidirectional stream for high-throughput sending. Each message gets an individual response confirming receipt.

upstream.go
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()
upstream.py
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()
upstream.ts
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();
Upstream.java
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();
Upstream.cs
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();
Upstream.kt
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()
upstream.cpp
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();
upstream.rs
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();
upstream.rb
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
upstream.exs
{: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.

downstream.go
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()
}
downstream.py
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,
)
downstream.ts
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));
Downstream.java
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());
Downstream.cs
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();
    }
}
Downstream.kt
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}") }
)
downstream.cpp
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;
    }
);
downstream.rs
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?;
downstream.rb
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
downstream.exs
# 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}")
end

Stream vs Polling Comparison

FeatureStream APIPolling (PollQueue)
ConnectionPersistent bidirectionalRequest/response per poll
LatencyLower (always connected)Higher (new request each time)
ThroughputHigherLower
Resource usageHolds connection openReleases between polls
Best forHigh-volume continuous processingPeriodic batch processing

Next Steps

Was this page helpful?

On this page