Stream Receive
Receive messages via streaming transport with manual ack
Overview
A downstream receiver is the persistent-connection way to pull queue messages: instead of opening and tearing down a request for every batch, you open one gRPC stream and reuse it across many receive cycles. That matters for any consumer that runs continuously — a worker loop, a background processor — where reconnecting per batch would add latency and churn on both the client and the broker.
receiveQueuesMessages uses that downstream stream internally, and with autoAck = false fetches a batch under manual settlement — nothing leaves the queue until you explicitly settle it. Each returned message is settled on its own: calling msg.ack() removes it from the queue immediately, while an unacknowledged message is redelivered once the visibility timeout expires.
Gotchas: a crash between receiving and acknowledging redelivers the whole batch, so processing must be idempotent; forgetting to call ack() doesn't lose the message, it just delays redelivery until waitTimeoutMs (or the visibility window) elapses; and leaving autoAck at its default silently removes messages on delivery, defeating the manual-settlement guarantee this pattern exists for.
Prerequisites
- KubeMQ server running on
localhost:50000 - Kotlin SDK installed (
implementation("io.kubemq.sdk:kubemq-sdk-kotlin:1.0.1"))
Code
package io.kubemq.sdk.examples.queuesstream
import io.kubemq.sdk.client.KubeMQClient
import io.kubemq.sdk.queues.queueMessage
import kotlinx.coroutines.runBlocking
private const val ADDRESS = "localhost:50000"
private const val CLIENT_ID = "kotlin-queues-stream-receive"
private const val CHANNEL = "kotlin-queues.stream-receive"
fun main() = runBlocking {
val client = KubeMQClient.queues {
address = ADDRESS
clientId = CLIENT_ID
}
client.use {
try {
client.createQueuesChannel(CHANNEL)
// Send messages to the queue
repeat(5) { i ->
client.sendQueuesMessage(queueMessage {
channel = CHANNEL
body = "Message ${i + 1}".toByteArray()
})
}
// Poll for messages via stream
println("Receiving messages via stream poll...\n")
val response = client.receiveQueuesMessages {
channel = CHANNEL
maxItems = 10
waitTimeoutMs = 5000
autoAck = false
}
// Process each message and acknowledge it
response.messages.forEach { msg ->
println(" Received: ${String(msg.body)}")
msg.ack()
}
println("\nReceived and acked ${response.messages.size} messages.")
} finally {
try { client.deleteQueuesChannel(CHANNEL) } catch (_: Exception) {}
}
println("Done.")
}
}How It Works
receiveQueuesMessages { }uses the gRPC downstream stream internally for efficient polling.- Each message is individually acknowledged with
msg.ack()after processing. - The
waitTimeoutMscontrols how long to wait for messages before returning. - Channel is created and cleaned up as part of the example.
Related
Was this page helpful?