Stream Receive
Poll KubeMQ queue messages via the downstream streaming API using the Elixir SDK.
Overview
A downstream receiver is the transactional way to pull queue messages: rather than a plain request/response fetch, it's backed by a gRPC stream and gives you an explicit settlement step before a batch is removed from the queue. That matters for any consumer that needs "fetch, process, then confirm" instead of "fetch and it's already gone."
KubeMQ.Client.poll_queue/2 fetches a batch under a transaction — nothing is removed from the queue until you explicitly settle it. Acknowledging the whole batch with KubeMQ.PollResponse.ack_all/1 permanently removes every message in one round-trip; until that call succeeds, the messages stay invisible to other consumers but are not yet gone.
Gotchas: a crash between polling and calling ack_all/1 redelivers the whole batch once the visibility timeout expires, so processing must be idempotent; ack_all/1 is all-or-nothing — a single bad message in the batch means you either ack everything (including messages you couldn't process) or nothing; and wait_timeout bounds how long the poll blocks on an empty queue, not how long processing may take afterward.
Prerequisites
- KubeMQ server running on
localhost:50000 - Elixir SDK installed (
{:kubemq, "~> 1.0"}in mix.exs)
Code
channel = "elixir-queues.stream-receive"
{:ok, client} = KubeMQ.Client.start_link(address: "localhost:50000", client_id: "elixir-qs-recv")
for i <- 1..3 do
{:ok, _} = KubeMQ.Client.send_queue_message(client,
KubeMQ.QueueMessage.new(channel: channel, body: "Queued item #{i}"))
end
IO.puts("Sent 3 messages")
case KubeMQ.Client.poll_queue(client,
channel: channel,
max_items: 10,
wait_timeout: 5_000
) do
{:ok, poll} ->
IO.puts("Polled #{length(poll.messages)} messages")
Enum.each(poll.messages, fn msg ->
IO.puts(" Channel: #{msg.channel}")
IO.puts(" Body: #{msg.body}")
end)
{:ok, _} = KubeMQ.PollResponse.ack_all(poll)
IO.puts("All messages acknowledged")
{:error, err} ->
IO.puts("Poll failed: #{err.message}")
end
KubeMQ.Client.close(client)How It Works
poll_queue/2uses the downstream streaming API to poll for messages with transactional control — distinct from the plainreceive_queue_messages/3pull, which has no settlement step- Messages must be explicitly acknowledged via
PollResponse.ack_all/1 - Until acknowledged, messages remain invisible to other consumers
Related
Was this page helpful?