KubeMQ
LearnEventsHow-To Guides

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.

channel_filter.go
// 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.>", "", ...)
channel_filter.py
# 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.>", ...))
channel_filter.js
// 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.>", ... });
ChannelFilter.java
// 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());
ChannelFilter.cs
// 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.*" })) { }
ChannelFilter.kt
// 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.*", ...)
channel_filter.cpp
// 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);
channel_filter.rs
// 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?;
channel_filter.rb
# 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) }
channel_filter.exs
# 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.

metadata_filter.go
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)
    }),
)
metadata_filter.py
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(),
)
metadata_filter.js
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),
});
MetadataFilter.java
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());
MetadataFilter.cs
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)}");
}
MetadataFilter.kt
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}") }
)
metadata_filter.cpp
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;
    }
);
metadata_filter.rs
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?;
metadata_filter.rb
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}"
end
metadata_filter.exs
KubeMQ.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

publish_with_tags.go
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",
    }),
)
publish_with_tags.py
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"},
    )
)
publish_with_tags.js
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" },
});
PublishWithTags.java
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());
PublishWithTags.cs
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"
    }
});
PublishWithTags.kt
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"),
))
publish_with_tags.cpp
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);
publish_with_tags.rs
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?;
publish_with_tags.rb
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" }
))
publish_with_tags.exs
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

tag_filter.go
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)
    }),
)
tag_filter.py
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(),
)
tag_filter.js
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),
});
TagFilter.java
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());
TagFilter.cs
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)}");
}
TagFilter.kt
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}") }
)
tag_filter.cpp
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;
    }
);
tag_filter.rs
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?;
tag_filter.rb
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}"
end
tag_filter.exs
KubeMQ.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

StrategyProsConsBest For
Channel patternsZero overhead, server-side filteringRequires channel naming disciplineHigh-throughput, predictable categories
Metadata filteringSimple, single field to checkAll events still delivered to subscriberEvent type discrimination
Tag filteringMulti-dimensional, flexibleAll events still delivered, parsing overheadComplex filtering criteria

For maximum efficiency, combine strategies: use channel hierarchy for coarse-grained filtering and tags for fine-grained filtering within a channel.

Was this page helpful?

On this page