Queues
Durable queues and streaming transports -- KubeMQ Kotlin SDK reference.
QueuesClient sends messages with sendQueuesMessage, receives batches with receiveQueuesMessages, and provides batch operations like ackAllQueuesMessages, nackAllQueuesMessages, and reQueueAllMessages. For non-destructive reads use peekQueueMessages with SimpleQueueReceiveConfig.
sendQueuesMessage
suspend fun sendQueuesMessage(message: QueueMessage): QueueSendResultSends a single queue message.
Parameters:
| Name | Type | Required | Description |
|---|---|---|---|
message | QueueMessage | Yes | Built via queueMessage { } DSL |
Returns: QueueSendResult -- message ID and error status.
sendQueueMessagesBatch
suspend fun sendQueueMessagesBatch(messages: List<QueueMessage>): List<QueueSendResult>Sends multiple messages in a single batch for high throughput.
receiveQueuesMessages
suspend fun receiveQueuesMessages(config: QueueReceiveConfig.() -> Unit): QueueReceiveResponseReceives batches of messages. Each QueueReceivedMessage exposes ack(), reject() unless autoAck is enabled.
Configuration DSL:
| Property | Type | Required | Description |
|---|---|---|---|
channel | String | Yes | Queue channel name |
maxItems | Int | Yes | Max messages to receive |
waitTimeoutMs | Int | Yes | Long poll wait (milliseconds) |
autoAck | Boolean | No | Acknowledge on delivery |
peekQueueMessages
suspend fun peekQueueMessages(config: SimpleQueueReceiveConfig.() -> Unit): SimpleQueueReceiveResponseNon-destructive read -- messages remain in the queue.
Configuration DSL (SimpleQueueReceiveConfig):
| Property | Type | Required | Description |
|---|---|---|---|
channel | String | Yes | Queue channel name |
maxNumberOfMessages | Int | No | Max messages to peek (default 1) |
waitTimeSeconds | Int | No | Poll wait timeout in seconds (default 1) |
Batch acknowledgement
suspend fun ackAllQueuesMessages(response: QueueReceiveResponse)
suspend fun nackAllQueuesMessages(response: QueueReceiveResponse)
suspend fun reQueueAllMessages(response: QueueReceiveResponse, channel: String)queueMessage DSL
queueMessage {
channel = "queues.tasks"
body = "process-this".toByteArray()
metadata = "task"
tags["priority"] = "high"
policy = QueueMessagePolicy(
delaySeconds = 5,
expirationSeconds = 60,
maxReceiveCount = 3,
maxReceiveQueue = "queues.dlq",
)
}| Property | Type | Description |
|---|---|---|
channel | String | Target queue |
body | ByteArray | Payload |
metadata | String | UTF-8 metadata |
tags | MutableMap<String, String> | Key/value tags |
policy | QueueMessagePolicy | Delay, expiration, DLQ config |
QueueMessagePolicy
| Field | Type | Description |
|---|---|---|
delaySeconds | Int | Initial delivery delay |
expirationSeconds | Int | Message TTL |
maxReceiveCount | Int | Max receive attempts before DLQ |
maxReceiveQueue | String | Dead letter queue channel |
Quick Usage
val queues = KubeMQClient.queues {
address = "localhost:50000"
clientId = "queues-demo"
}
queues.use {
queues.sendQueuesMessage(queueMessage {
channel = "jobs"
body = payload.toByteArray()
})
}See Also
Was this page helpful?