How to Implement Topic Exchanges with Wildcard Routing in Kourier
Use BuiltinExchangeType.TOPIC when declaring the exchange, then bind queues with routing keys containing * (match one word) or # (match zero or more words) to enable pattern-based message routing.
Kourier is a Kotlin Multiplatform implementation of AMQP v0-9-1 that provides type-safe DSL builders for exchange and queue operations. Topic exchanges allow publishers to route messages based on hierarchical routing keys, while consumers use wildcard binding keys to receive selective message streams. This guide demonstrates how to declare topic exchanges, publish with routing keys, and bind queues using wildcard patterns in the Kourier client.
Understanding Topic Exchanges and Wildcard Routing
A topic exchange routes messages to queues based on wildcard matching between the message routing key and the queue binding key. In BuiltinExchangeType.kt, the topic exchange is defined as a standard built-in type:
// amqp-core/src/commonMain/kotlin/dev/kourier/amqp/BuiltinExchangeType.kt
enum class BuiltinExchangeType(val type: String) {
DIRECT("direct"),
FANOUT("fanout"),
TOPIC("topic"),
HEADERS("headers")
}
Routing keys consist of dot-separated words (e.g., stock.usd.nyse, kern.critical). Binding keys support two wildcards:
*— Matches exactly one word in a routing key#— Matches zero or more words in a routing key
When a message arrives, the broker (RabbitMQ) evaluates the routing key against all binding keys; matching queues receive the message. The binding operation uses the routingKey field in the AMQP Exchange.Bind frame defined in Frame.kt.
Declaring a Topic Exchange in Kourier
Use the exchangeDeclare extension function on AMQPChannel to create a topic exchange. This function accepts a DeclaredExchangeBuilder block that maps to the protocol Exchange.Declare frame.
// amqp-client/src/commonMain/kotlin/dev/kourier/amqp/channel/Extensions.kt
suspend fun AMQPChannel.exchangeDeclare(block: DeclaredExchangeBuilder.() -> Unit)
The builder properties are defined in DeclaredExchangeBuilder.kt:
// amqp-client/src/commonMain/kotlin/dev/kourier/amqp/states/DeclaredExchangeBuilder.kt
class DeclaredExchangeBuilder {
var name: String = ""
var type: BuiltinExchangeType = BuiltinExchangeType.DIRECT
var durable: Boolean = false
var autoDelete: Boolean = false
var internal: Boolean = false
var arguments: Map<String, Any?> = emptyMap()
}
To declare a topic exchange, set type to BuiltinExchangeType.TOPIC:
channel.exchangeDeclare {
name = "topic_logs"
type = BuiltinExchangeType.TOPIC
durable = false
autoDelete = false
internal = false
arguments = emptyMap()
}
Publishing Messages with Routing Keys
After declaring the exchange, publish messages using basicPublish with a hierarchical routing key. The routing key must follow the dot-notation format to enable wildcard matching.
val routingKey = "kern.critical" // facility.severity format
val message = "Kernel panic detected"
channel.basicPublish(
body = message.toByteArray(),
exchange = "topic_logs",
routingKey = routingKey,
properties = Properties()
)
The broker stores the routing key with the message and evaluates it against all queue binding keys when determining delivery targets.
Binding Queues with Wildcard Patterns
Bind queues to the topic exchange using queueBind, providing a binding key that contains wildcards. This creates an Exchange.Bind frame where the routingKey field holds the binding pattern.
// amqp-core/src/commonMain/kotlin/dev/kourier/amqp/Frame.kt
// Frame.Exchange.Bind carries the binding key in the routingKey field
Use the queueBind extension to establish the relationship:
channel.queueBind(
queue = queueName,
exchange = "topic_logs",
routingKey = "kern.*" // Matches kern.critical, kern.info, etc.
)
Common wildcard patterns include:
#— Receive all messages regardless of routing keykern.*— Receive all messages from thekernfacility (any severity)*.critical— Receive all critical severity messages (any facility)kern.critical— Receive onlykern.criticalmessages (exact match)
A single queue can bind multiple patterns to aggregate different message streams:
listOf("kern.*", "*.critical").forEach { pattern ->
channel.queueBind(queueName, "topic_logs", pattern)
}
Complete Producer and Consumer Example
The following implementation demonstrates a complete topic exchange workflow using Kourier's coroutine-based API.
Producer Implementation
suspend fun emitLogTopic(
coroutineScope: CoroutineScope,
routingKey: String,
message: String
) {
val config = amqpConfig {
server { host = "localhost" }
}
val connection = createAMQPConnection(coroutineScope, config)
val channel = connection.openChannel()
// Declare a topic exchange
channel.exchangeDeclare(
"topic_logs",
BuiltinExchangeType.TOPIC,
durable = false,
autoDelete = false,
internal = false,
arguments = emptyMap()
)
// Publish with routing key
channel.basicPublish(
body = message.toByteArray(),
exchange = "topic_logs",
routingKey = routingKey,
properties = Properties()
)
println(" [x] Sent '$routingKey':'$message'")
channel.close()
connection.close()
}
Consumer Implementation
suspend fun receiveLogsTopic(
coroutineScope: CoroutineScope,
bindingKeys: List<String>
) {
val config = amqpConfig {
server { host = "localhost" }
}
val connection = createAMQPConnection(coroutineScope, config)
val channel = connection.openChannel()
// Declare the topic exchange
channel.exchangeDeclare(
"topic_logs",
BuiltinExchangeType.TOPIC,
durable = false,
autoDelete = false,
internal = false,
arguments = emptyMap()
)
// Create a temporary exclusive queue
val queueDeclared = channel.queueDeclare(
name = "",
durable = false,
exclusive = true,
autoDelete = true,
arguments = emptyMap()
)
val queueName = queueDeclared.queueName
// Bind with wildcard patterns
for (bindingKey in bindingKeys) {
channel.queueBind(
queue = queueName,
exchange = "topic_logs",
routingKey = bindingKey
)
}
println(" [*] Waiting for logs. To exit press CTRL+C")
// Consume messages
val consumer = channel.basicConsume(queueName, noAck = true)
for (delivery in consumer) {
val rk = delivery.message.routingKey
val msg = delivery.message.body.decodeToString()
println(" [x] Received '$rk':'$msg'")
}
channel.close()
connection.close()
}
Usage Examples
Receive all messages using the hash wildcard:
receiveLogsTopic(this, listOf("#"))
Receive only kernel facility logs:
receiveLogsTopic(this, listOf("kern.*"))
Receive only critical severity from any facility:
receiveLogsTopic(this, listOf("*.critical"))
Combine multiple patterns to aggregate specific streams:
receiveLogsTopic(this, listOf("kern.*", "*.critical"))
Publish a message with a specific routing key:
emitLogTopic(this, "kern.critical", "A critical kernel error")
Summary
- Topic exchanges in Kourier use
BuiltinExchangeType.TOPICdefined inBuiltinExchangeType.ktto enable pattern-based routing. - Routing keys follow dot-notation (e.g.,
facility.severity) and are evaluated against binding keys containing wildcards. - Wildcard characters:
*matches exactly one word, while#matches zero or more words. - Declaration uses
exchangeDeclarewith aDeclaredExchangeBuilderblock to configure the exchange properties. - Binding uses
queueBindwith theroutingKeyparameter set to the wildcard pattern; this creates anExchange.Bindframe as defined inFrame.kt. - Publishing uses
basicPublishwith a hierarchical routing key; the broker handles wildcard matching server-side.
Frequently Asked Questions
What is the difference between * and # wildcards in Kourier topic exchanges?
The * (asterisk) wildcard matches exactly one word in a routing key, while the # (hash) wildcard matches zero or more words. For example, binding key kern.* matches kern.critical but not kern.critical.error, whereas kern.# matches both. This follows the AMQP v0-9-1 specification as implemented in the Kourier client.
How do I declare a durable topic exchange in Kourier?
Set the durable property to true in the DeclaredExchangeBuilder block when calling exchangeDeclare. This persists the exchange metadata to disk so it survives broker restarts. The builder class in DeclaredExchangeBuilder.kt captures this flag and maps it to the durable bit in the AMQP Exchange.Declare frame.
channel.exchangeDeclare {
name = "durable_topic"
type = BuiltinExchangeType.TOPIC
durable = true
}
Can a queue bind to multiple routing patterns in Kourier?
Yes, a single queue can bind to a topic exchange multiple times with different binding keys. Invoke queueBind repeatedly with distinct routingKey values (e.g., "kern.*" and "*.critical"). Each call creates a separate Exchange.Bind frame as defined in Frame.kt, allowing the queue to aggregate messages matching any of the specified patterns.
Where does the wildcard matching logic execute in Kourier?
The wildcard matching logic executes on the broker side (e.g., RabbitMQ), not within the Kourier client. Kourier's responsibility is to transmit the routing key and binding key values via the AMQP protocol frames. The Frame.Exchange.Bind class in Frame.kt carries the binding key in its routingKey field, which the broker evaluates against incoming message routing keys using the * and # wildcard rules.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →