KubeMQ
Client SDKsKotlinReference

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): QueueSendResult

Sends a single queue message.

Parameters:

NameTypeRequiredDescription
messageQueueMessageYesBuilt 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): QueueReceiveResponse

Receives batches of messages. Each QueueReceivedMessage exposes ack(), reject() unless autoAck is enabled.

Configuration DSL:

PropertyTypeRequiredDescription
channelStringYesQueue channel name
maxItemsIntYesMax messages to receive
waitTimeoutMsIntYesLong poll wait (milliseconds)
autoAckBooleanNoAcknowledge on delivery

peekQueueMessages

suspend fun peekQueueMessages(config: SimpleQueueReceiveConfig.() -> Unit): SimpleQueueReceiveResponse

Non-destructive read -- messages remain in the queue.

Configuration DSL (SimpleQueueReceiveConfig):

PropertyTypeRequiredDescription
channelStringYesQueue channel name
maxNumberOfMessagesIntNoMax messages to peek (default 1)
waitTimeSecondsIntNoPoll 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",
    )
}
PropertyTypeDescription
channelStringTarget queue
bodyByteArrayPayload
metadataStringUTF-8 metadata
tagsMutableMap<String, String>Key/value tags
policyQueueMessagePolicyDelay, expiration, DLQ config

QueueMessagePolicy

FieldTypeDescription
delaySecondsIntInitial delivery delay
expirationSecondsIntMessage TTL
maxReceiveCountIntMax receive attempts before DLQ
maxReceiveQueueStringDead letter queue channel

Quick Usage

Queues.kt
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?

On this page