Filter Events by Tags
Use message tags and metadata to filter events on the subscriber side.
KubeMQ Events does not have a built-in server-side filter mechanism, but you can implement effective filtering using three strategies: channel hierarchy, metadata inspection, and tag-based routing.
The diagram below shows the tag-filter decision: every subscriber receives all events on the channel, then applies its own predicate — accepting matching events and ignoring the rest.
Tag filtering happens in the subscriber: KubeMQ delivers every event, the client keeps only those whose tags match.
Strategy 1: Filter by Channel Pattern
The most efficient approach — organize events into hierarchical channels and use wildcard subscriptions to select the subset you need.
// Publish to specific channels
client.SendEvent(ctx, kubemq.NewEvent().
SetChannel("orders.us-east.created").
SetBody([]byte(orderData)))
// Subscribe to a filtered subset
client.SubscribeToEvents(ctx, "orders.us-east.*", "", ...)
client.SubscribeToEvents(ctx, "orders.*.created", "", ...)
client.SubscribeToEvents(ctx, "orders.>", "", ...)# Publish to specific channels
client.send_event(EventMessage(
channel="orders.us-east.created", body=order_data
))
# Subscribe to a filtered subset
client.subscribe_to_events(EventsSubscription(channel="orders.us-east.*", ...))
client.subscribe_to_events(EventsSubscription(channel="orders.*.created", ...))
client.subscribe_to_events(EventsSubscription(channel="orders.>", ...))// Publish to specific channels
await client.sendEvent({ channel: "orders.us-east.created", body: orderData });
// Subscribe to a filtered subset
client.subscribeToEvents({ channel: "orders.us-east.*", ... });
client.subscribeToEvents({ channel: "orders.*.created", ... });
client.subscribeToEvents({ channel: "orders.>", ... });// Publish to specific channels
client.sendEventsMessage(EventMessage.builder()
.channel("orders.us-east.created")
.body(orderData).build());
// Subscribe to a filtered subset
client.subscribeToEvents(EventsSubscription.builder()
.channel("orders.us-east.*").build());// Publish to specific channels
await client.SendEventAsync(new EventMessage
{
Channel = "orders.us-east.created",
Body = orderData
});
// Subscribe to a filtered subset
await foreach (var msg in client.SubscribeToEventsAsync(
new EventsSubscription { Channel = "orders.us-east.*" })) { }// Publish to specific channels
client.sendEvent(EventMessage(
channel = "orders.us-east.created", body = orderData
))
// Subscribe to a filtered subset
client.subscribeToEvents(channel = "orders.us-east.*", ...)// Publish to specific channels
event.channel = "orders.us-east.created";
client.sendEvent(event);
// Subscribe to a filtered subset
client.subscribeToEvents("orders.us-east.*", "", handler, errHandler);// Publish to specific channels
let event = EventBuilder::new()
.channel("orders.us-east.created")
.body(order_data)
.build();
client.send_event(event).await?;
// Subscribe to a filtered subset
client.subscribe_to_events("orders.us-east.*", "", handler, None).await?;
client.subscribe_to_events("orders.*.created", "", handler, None).await?;
client.subscribe_to_events("orders.>", "", handler, None).await?;# Publish to specific channels
client.send_event(KubeMQ::PubSub::EventMessage.new(
channel: "orders.us-east.created", body: order_data
))
# Subscribe to a filtered subset
client.subscribe_to_events(
KubeMQ::PubSub::EventsSubscription.new(channel: "orders.us-east.*"),
cancellation_token: cancel
) { |event| handle(event) }
client.subscribe_to_events(
KubeMQ::PubSub::EventsSubscription.new(channel: "orders.>"),
cancellation_token: cancel
) { |event| handle(event) }# Publish to specific channels
event = KubeMQ.Event.new(channel: "orders.us-east.created", body: order_data)
KubeMQ.Client.send_event(client, event)
# Subscribe to a filtered subset
KubeMQ.Client.subscribe_to_events(client, "orders.us-east.*", on_event: &handle/1)
KubeMQ.Client.subscribe_to_events(client, "orders.>", on_event: &handle/1)This approach has zero runtime overhead because filtering happens at the broker subscription level before messages reach your code.
Strategy 2: Filter by Metadata
Use the metadata field to carry event type information and filter in your subscriber callback.
sub, err := client.SubscribeToEvents(ctx, "order-events", "",
kubemq.WithOnEvent(func(event *kubemq.Event) {
if event.Metadata != "order.created" {
return
}
fmt.Printf("Processing new order: %s\n", string(event.Body))
}),
kubemq.WithOnError(func(err error) {
log.Println("Error:", err)
}),
)def on_event(event):
if event.metadata != "order.created":
return
print(f"Processing new order: {event.body.decode('utf-8')}")
client.subscribe_to_events(
subscription=EventsSubscription(
channel="order-events",
on_receive_event_callback=on_event,
on_error_callback=lambda e: print(f"Error: {e}"),
),
cancel=CancellationToken(),
)client.subscribeToEvents({
channel: "order-events",
onEvent: (msg) => {
if (msg.metadata !== "order.created") return;
console.log(`Processing new order: ${Buffer.from(msg.body).toString()}`);
},
onError: (err) => console.error("Error:", err.message),
});client.subscribeToEvents(EventsSubscription.builder()
.channel("order-events")
.onReceiveEventCallback(event -> {
if (!"order.created".equals(event.getMetadata())) return;
System.out.println("Processing new order: "
+ new String(event.getBody()));
})
.onErrorCallback(err ->
System.err.println("Error: " + err.getMessage()))
.build());await foreach (var msg in client.SubscribeToEventsAsync(
new EventsSubscription { Channel = "order-events" }))
{
if (msg.Metadata != "order.created") continue;
Console.WriteLine($"Processing new order: "
+ $"{Encoding.UTF8.GetString(msg.Body.Span)}");
}client.subscribeToEvents(
channel = "order-events",
onEvent = { event ->
if (event.metadata != "order.created") return@subscribeToEvents
println("Processing new order: ${String(event.body)}")
},
onError = { err -> System.err.println("Error: ${err.message}") }
)client.subscribeToEvents("order-events", "",
[](const kubemq::Event& event) {
if (event.metadata != "order.created") return;
std::cout << "Processing new order: " << event.body << std::endl;
},
[](const std::string& err) {
std::cerr << "Error: " << err << std::endl;
}
);let sub = client
.subscribe_to_events(
"order-events",
"",
|event| {
Box::pin(async move {
if event.metadata != "order.created" {
return;
}
println!(
"Processing new order: {}",
String::from_utf8_lossy(&event.body)
);
})
},
None,
)
.await?;sub = KubeMQ::PubSub::EventsSubscription.new(channel: "order-events")
client.subscribe_to_events(
sub,
cancellation_token: cancel,
on_error: ->(e) { puts "Error: #{e.message}" }
) do |event|
next unless event.metadata == "order.created"
puts "Processing new order: #{event.body}"
endKubeMQ.Client.subscribe_to_events(client, "order-events",
on_event: fn event ->
if event.metadata == "order.created" do
IO.puts("Processing new order: #{event.body}")
end
end,
on_error: fn err -> IO.puts("Error: #{err.message}") end
)Client-side metadata filtering still delivers all events to the subscriber. Use channel-based filtering (Strategy 1) when throughput is high and you want to reduce network traffic.
Strategy 3: Filter by Tags
Tags are key-value pairs attached to events. Use them for multi-dimensional filtering when a single channel hierarchy is not sufficient.
Publish with Tags
err = client.SendEvent(ctx, kubemq.NewEvent().
SetChannel("order-events").
SetBody([]byte(`{"orderId":"ORD-100","amount":250.00}`)).
SetMetadata("order.created").
SetTags(map[string]string{
"region": "us-east",
"priority": "high",
"customer": "enterprise",
}),
)client.send_event(
EventMessage(
channel="order-events",
body=b'{"orderId":"ORD-100","amount":250.00}',
metadata="order.created",
tags={"region": "us-east", "priority": "high", "customer": "enterprise"},
)
)await client.sendEvent({
channel: "order-events",
body: Buffer.from('{"orderId":"ORD-100","amount":250.00}'),
metadata: "order.created",
tags: { region: "us-east", priority: "high", customer: "enterprise" },
});client.sendEventsMessage(EventMessage.builder()
.channel("order-events")
.body("{\"orderId\":\"ORD-100\",\"amount\":250.00}".getBytes())
.metadata("order.created")
.tags(Map.of("region", "us-east", "priority", "high", "customer", "enterprise"))
.build());await client.SendEventAsync(new EventMessage
{
Channel = "order-events",
Body = Encoding.UTF8.GetBytes("{\"orderId\":\"ORD-100\",\"amount\":250.00}"),
Metadata = "order.created",
Tags = new Dictionary<string, string>
{
["region"] = "us-east",
["priority"] = "high",
["customer"] = "enterprise"
}
});client.sendEvent(EventMessage(
channel = "order-events",
body = """{"orderId":"ORD-100","amount":250.00}""".toByteArray(),
metadata = "order.created",
tags = mapOf("region" to "us-east", "priority" to "high", "customer" to "enterprise"),
))kubemq::EventMessage event;
event.channel = "order-events";
event.body = R"({"orderId":"ORD-100","amount":250.00})";
event.metadata = "order.created";
event.tags["region"] = "us-east";
event.tags["priority"] = "high";
event.tags["customer"] = "enterprise";
client.sendEvent(event);use std::collections::HashMap;
let tags = HashMap::from([
("region".to_string(), "us-east".to_string()),
("priority".to_string(), "high".to_string()),
("customer".to_string(), "enterprise".to_string()),
]);
let event = EventBuilder::new()
.channel("order-events")
.metadata("order.created")
.tags(tags)
.body(br#"{"orderId":"ORD-100","amount":250.00}"#.to_vec())
.build();
client.send_event(event).await?;client.send_event(KubeMQ::PubSub::EventMessage.new(
channel: "order-events",
body: '{"orderId":"ORD-100","amount":250.00}',
metadata: "order.created",
tags: { "region" => "us-east", "priority" => "high", "customer" => "enterprise" }
))event =
KubeMQ.Event.new(
channel: "order-events",
body: ~s({"orderId":"ORD-100","amount":250.00}),
metadata: "order.created",
tags: %{"region" => "us-east", "priority" => "high", "customer" => "enterprise"}
)
KubeMQ.Client.send_event(client, event)Filter by Tags in the Subscriber
sub, err := client.SubscribeToEvents(ctx, "order-events", "",
kubemq.WithOnEvent(func(event *kubemq.Event) {
if event.Tags["priority"] != "high" ||
event.Tags["customer"] != "enterprise" {
return
}
fmt.Printf("High-priority enterprise order: %s\n",
string(event.Body))
}),
kubemq.WithOnError(func(err error) {
log.Println("Error:", err)
}),
)def on_event(event):
if event.tags.get("priority") != "high" or \
event.tags.get("customer") != "enterprise":
return
print(f"High-priority enterprise order: {event.body.decode('utf-8')}")
client.subscribe_to_events(
subscription=EventsSubscription(
channel="order-events",
on_receive_event_callback=on_event,
on_error_callback=lambda e: print(f"Error: {e}"),
),
cancel=CancellationToken(),
)client.subscribeToEvents({
channel: "order-events",
onEvent: (msg) => {
if (msg.tags?.priority !== "high" || msg.tags?.customer !== "enterprise") {
return;
}
console.log(
`High-priority enterprise order: ${Buffer.from(msg.body).toString()}`
);
},
onError: (err) => console.error("Error:", err.message),
});client.subscribeToEvents(EventsSubscription.builder()
.channel("order-events")
.onReceiveEventCallback(event -> {
Map<String, String> tags = event.getTags();
if (!"high".equals(tags.get("priority")) ||
!"enterprise".equals(tags.get("customer"))) return;
System.out.println("High-priority enterprise order: "
+ new String(event.getBody()));
})
.onErrorCallback(err ->
System.err.println("Error: " + err.getMessage()))
.build());await foreach (var msg in client.SubscribeToEventsAsync(
new EventsSubscription { Channel = "order-events" }))
{
if (msg.Tags?.GetValueOrDefault("priority") != "high" ||
msg.Tags?.GetValueOrDefault("customer") != "enterprise") continue;
Console.WriteLine($"High-priority enterprise order: "
+ $"{Encoding.UTF8.GetString(msg.Body.Span)}");
}client.subscribeToEvents(
channel = "order-events",
onEvent = { event ->
if (event.tags["priority"] != "high" ||
event.tags["customer"] != "enterprise") return@subscribeToEvents
println("High-priority enterprise order: ${String(event.body)}")
},
onError = { err -> System.err.println("Error: ${err.message}") }
)client.subscribeToEvents("order-events", "",
[](const kubemq::Event& event) {
if (event.tags.at("priority") != "high" ||
event.tags.at("customer") != "enterprise") return;
std::cout << "High-priority enterprise order: "
<< event.body << std::endl;
},
[](const std::string& err) {
std::cerr << "Error: " << err << std::endl;
}
);let sub = client
.subscribe_to_events(
"order-events",
"",
|event| {
Box::pin(async move {
let high = event.tags.get("priority").map(String::as_str) == Some("high");
let ent = event.tags.get("customer").map(String::as_str) == Some("enterprise");
if !high || !ent {
return;
}
println!(
"High-priority enterprise order: {}",
String::from_utf8_lossy(&event.body)
);
})
},
None,
)
.await?;sub = KubeMQ::PubSub::EventsSubscription.new(channel: "order-events")
client.subscribe_to_events(
sub,
cancellation_token: cancel,
on_error: ->(e) { puts "Error: #{e.message}" }
) do |event|
next unless event.tags["priority"] == "high" &&
event.tags["customer"] == "enterprise"
puts "High-priority enterprise order: #{event.body}"
endKubeMQ.Client.subscribe_to_events(client, "order-events",
on_event: fn event ->
if event.tags["priority"] == "high" and
event.tags["customer"] == "enterprise" do
IO.puts("High-priority enterprise order: #{event.body}")
end
end,
on_error: fn err -> IO.puts("Error: #{err.message}") end
)Choosing a Strategy
| Strategy | Pros | Cons | Best For |
|---|---|---|---|
| Channel patterns | Zero overhead, server-side filtering | Requires channel naming discipline | High-throughput, predictable categories |
| Metadata filtering | Simple, single field to check | All events still delivered to subscriber | Event type discrimination |
| Tag filtering | Multi-dimensional, flexible | All events still delivered, parsing overhead | Complex filtering criteria |
For maximum efficiency, combine strategies: use channel hierarchy for coarse-grained filtering and tags for fine-grained filtering within a channel.
Related
- Wildcard Subscriptions for channel-based filtering patterns
- Multicast Events for routing events across channels
- Events Reference for tag format and validation rules
Was this page helpful?