Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Message filtering

Última atualização: Jun 27, 2026

A filtragem de mensagens permite que o broker do ApsaraMQ for RocketMQ entregue aos consumidores apenas as mensagens correspondentes a condições específicas. Os produtores classificam as mensagens ao definir atributos, e os consumidores especificam condições de filtro durante a assinatura. O broker avalia cada mensagem conforme essas condições e entrega somente aquelas que atendem aos critérios.

Se um consumidor assinar um tópico sem definir uma condição de filtro, receberá todas as mensagens desse tópico, independentemente dos atributos configurados.

Métodos de filtragem

O ApsaraMQ for RocketMQ oferece suporte a dois métodos de filtragem.

Método

Descrição

Quando usar

Requisito da instância

Requisito de protocolo

Filtragem baseada em tag (padrão)

Atribua uma tag a cada mensagem no lado do produtor. Especifique quais tags receber no lado do consumidor. O broker entrega as mensagens cuja tag corresponde à assinatura.

Para filtragem por um único atributo. Apenas uma tag por mensagem.

Nenhum

Nenhum

Filtragem SQL baseada em atributos

Defina vários atributos personalizados em cada mensagem. Crie expressões SQL no lado do consumidor para filtrar pelos valores desses atributos.

Para filtragem por múltiplos atributos com condições complexas.

Somente Enterprise Platinum Edition

Somente SDK cliente TCP

Filtragem baseada em tag

Uma tag é um rótulo usado para classificar mensagens dentro de um tópico. O produtor atribui uma tag a cada mensagem antes do envio, e o consumidor assina as mensagens com base em tags específicas.

Cenário de exemplo

Em um sistema de transações de e-commerce, há três tipos de mensagens:

  • Mensagens de pedido

  • Mensagens de pagamento

  • Mensagens de logística

Essas mensagens são enviadas para um tópico chamado Trade_Topic. Quatro sistemas downstream assinam este tópico:

  • Sistema de pagamentos: assina apenas mensagens de pagamento.

  • Sistema de logística: assina apenas mensagens de logística.

  • Sistema de análise de taxa de sucesso de transações: assina mensagens de pedido e de pagamento.

  • Sistema de computação em tempo real: assina todas as mensagens.

Tag-based filtering process

Configurar filtragem baseada em tag

Defina a lógica de filtragem no SDK cliente do ApsaraMQ for RocketMQ. Configure tags nas mensagens antes do envio e especifique quais tags assinar no lado do consumidor. Para mais detalhes sobre o SDK, consulte Visão geral.

Enviar uma mensagem com tag

Especifique uma tag para cada mensagem antes de enviar:

Message msg = new Message("MQ_TOPIC","TagA","Hello MQ".getBytes());

Assinar todas as mensagens

Use um asterisco (*) para assinar todas as mensagens de um tópico:

consumer.subscribe("MQ_TOPIC", "*", new MessageListener() {
        public Action consume(Message message, ConsumeContext context) {
            System.out.println(message.getMsgID());
            return Action.CommitMessage;
        }
    });

Assinar mensagens com uma tag específica

Especifique a tag para receber apenas as mensagens correspondentes:

consumer.subscribe("MQ_TOPIC", "TagA", new MessageListener() {
        public Action consume(Message message, ConsumeContext context) {
            System.out.println(message.getMsgID());
            return Action.CommitMessage;
        }
    });

Assinar mensagens com múltiplas tags

Separe várias tags com duas barras verticais (||):

consumer.subscribe("MQ_TOPIC", "TagA||TagB", new MessageListener() {
        public Action consume(Message message, ConsumeContext context) {
            System.out.println(message.getMsgID());
            return Action.CommitMessage;
        }
    });

Incorreto: múltiplas assinaturas para o mesmo tópico

Se um consumidor chamar consumer.subscribe() várias vezes para o mesmo tópico com tags diferentes, apenas a última assinatura terá efeito. As assinaturas anteriores serão substituídas.

// The consumer receives only messages with TagB. Messages with TagA are NOT received.
consumer.subscribe("MQ_TOPIC", "TagA", new MessageListener() {
        public Action consume(Message message, ConsumeContext context) {
            System.out.println(message.getMsgID());
            return Action.CommitMessage;
        }
    });
consumer.subscribe("MQ_TOPIC", "TagB", new MessageListener() {
        public Action consume(Message message, ConsumeContext context) {
            System.out.println(message.getMsgID());
            return Action.CommitMessage;
        }
    });

Filtragem SQL baseada em atributos

A filtragem SQL baseada em atributos permite que os produtores definam vários atributos personalizados em cada mensagem. Os consumidores criam expressões SQL para filtrar as mensagens com base nesses atributos. O broker avalia cada mensagem em relação à expressão SQL e entrega apenas aquelas que correspondem aos critérios.

Uma tag é um tipo especial de atributo de mensagem. A filtragem SQL baseada em atributos é compatível com a filtragem baseada em tag. Na sintaxe SQL, referencie o atributo de tag como TAGS .

Limitações

  • Somente instâncias Enterprise Platinum Edition oferecem suporte à filtragem SQL baseada em atributos.

  • Apenas SDKs clientes TCP suportam este método.

  • Caso o broker não ofereça suporte à filtragem SQL baseada em atributos e um consumidor defina uma expressão de filtro, ocorrerá um erro na inicialização do consumidor ou ele poderá deixar de receber mensagens.

Cenário de exemplo

Em um sistema de transações de e-commerce, as mensagens são classificadas em mensagens de pedido e mensagens de logística. As mensagens de logística incluem um atributo region com valores Hangzhou ou Shanghai.

  • Mensagens de pedido

  • Mensagens de logística

    • Mensagens de logística com region = Hangzhou

    • Mensagens de logística com region = Shanghai

Essas mensagens são enviadas para um tópico chamado Trade_Topic. Quatro sistemas downstream realizam a assinatura:

  • Sistema de logística 1: assina apenas mensagens de logística onde region é Hangzhou.

  • Sistema de logística 2: assina todas as mensagens de logística.

  • Sistema de rastreamento de pedidos: assina apenas mensagens de pedido.

  • Sistema de computação em tempo real: assina todas as mensagens.

Attribute-based SQL filtering process

Configurar filtragem SQL

Defina a lógica de filtragem no SDK cliente do ApsaraMQ for RocketMQ. Configure atributos personalizados de mensagem no código do produtor e defina uma expressão de filtro SQL no código do consumidor. Para mais detalhes sobre o SDK, consulte Visão geral.

Regras para chaves de atributos

Cada atributo personalizado é um par chave-valor. A chave deve seguir estas regras:

  • Pode conter letras, dígitos e sublinhados (_).

  • Deve começar com uma letra ou um sublinhado (_).

É possível definir múltiplos atributos em cada mensagem.

Produtor: definir atributos personalizados

Message msg = new Message("topic", "tagA", "Hello MQ".getBytes());
// Set custom attribute A to 1.
msg.putUserProperties("A", "1");

Consumidor: definir uma expressão de filtro

Utilize MessageSelector.bySql() para definir uma expressão de filtro SQL.

Importante

Sempre inclua uma verificação IS NOT NULL para cada atributo personalizado na expressão de filtro. Se o atributo não existir em uma mensagem, a expressão resultará em NULL e a mensagem não será entregue.

// Subscribe to messages where attribute A exists and equals 1.
consumer.subscribe("topic", MessageSelector.bySql("A IS NOT NULL AND TAGS IS NOT NULL AND A = '1'"), new MessageListener() {
    public Action consume(Message message, ConsumeContext context) {
        System.out.printf("Receive New Messages: %s %n", message);
        return Action.CommitMessage;
    }
});

Referência de sintaxe SQL

Sintaxe

Descrição

Exemplo

IS NULL

O atributo não existe.

a IS NULL -- O atributo a não existe.

IS NOT NULL

O atributo existe.

a IS NOT NULL -- O atributo a existe.

>, >=, <, <=

Comparação numérica. Não compara strings. Um erro será reportado na inicialização do consumidor se usado com strings. Strings conversíveis em números são tratadas como valores numéricos.

a IS NOT NULL AND a > 100 -- O atributo a existe e é maior que 100. a IS NOT NULL AND a > 'abc' -- Erro. Não é possível comparar a com a string abc.

BETWEEN xxx AND xxx

Faixa numérica, inclusiva. Equivalente a >= xxx AND <= xxx. Não compara strings.

a IS NOT NULL AND (a BETWEEN 10 AND 100) -- O atributo a existe e está entre 10 e 100, inclusive.

NOT BETWEEN xxx AND xxx

Fora de uma faixa numérica. Equivalente a < xxx OR > xxx. Não compara strings.

a IS NOT NULL AND (a NOT BETWEEN 10 AND 100) -- O atributo a existe e é menor que 10 ou maior que 100.

IN (xxx, xxx)

O valor está em um conjunto. Os elementos do conjunto devem ser strings.

a IS NOT NULL AND (a IN ('abc', 'def')) -- O atributo a existe e é igual a abc ou def.

=, <>

Igual e diferente de. Funciona tanto com números quanto com strings.

a IS NOT NULL AND (a = 'abc' OR a<>'def') -- O atributo a existe e é igual a abc ou diferente de def.

AND, OR

Operadores lógicos. Coloque cada condição entre parênteses.

a IS NOT NULL AND (a > 100) OR (b IS NULL) -- O atributo a existe e é maior que 100, ou o atributo b não existe.

Tratamento de erros em expressões de filtro

Quando uma expressão de filtro não produz um resultado válido, o broker não entrega a mensagem. Isso ocorre nas seguintes situações:

  • Erro de cálculo: Ocorre uma exceção durante a avaliação da expressão, como ao comparar valores numéricos e não numéricos.

  • Resultado NULL ou não booleano: A expressão retorna NULL ou um valor não booleano. Por exemplo, quando o consumidor filtra por um atributo que o produtor não definiu.

  • Incompatibilidade de tipo: O valor do atributo personalizado é um número de ponto flutuante, mas a expressão de filtro usa um inteiro para comparação.

Código de exemplo

Estes exemplos utilizam uma única mensagem de produtor e múltiplas assinaturas de consumidor para demonstrar diferentes resultados de filtragem.

Enviar uma mensagem com tag e atributos personalizados

Producer producer = ONSFactory.createProducer(properties);
// Set the tag to tagA.
Message msg = new Message("topicA", "tagA", "Hello MQ".getBytes());
// Set custom attribute region to hangzhou.
msg.putUserProperties("region", "hangzhou");
// Set custom attribute price to 50.
msg.putUserProperties("price", "50");
SendResult sendResult = producer.send(msg);

Filtrar por um atributo personalizado

Consumer consumer = ONSFactory.createConsumer(properties);
// Subscribe only to messages where region is hangzhou.
// Messages without the region attribute or with a different value are not delivered.
consumer.subscribe("topicA", MessageSelector.bySql("region IS NOT NULL AND region = 'hangzhou'"), new MessageListener() {
    public Action consume(Message message, ConsumeContext context) {
        System.out.printf("Receive New Messages: %s %n", message);
        return Action.CommitMessage;
    }
});

Resultado esperado: A mensagem é entregue. Ela possui o atributo region definido como hangzhou, o que corresponde à condição de filtro.

Filtrar por tag e atributo personalizado

Consumer consumer = ONSFactory.createConsumer(properties);
// Subscribe only to messages with tagA and price greater than 30.
consumer.subscribe("topicA", MessageSelector.bySql("TAGS IS NOT NULL AND price IS NOT NULL AND TAGS = 'tagA' AND price > 30 "), new MessageListener() {
    public Action consume(Message message, ConsumeContext context) {
        System.out.printf("Receive New Messages: %s %n", message);
        return Action.CommitMessage;
    }
});

Resultado esperado: A mensagem é entregue. Ela possui tagA e o atributo price é 50, valor maior que 30.

Filtrar por múltiplos atributos personalizados

Consumer consumer = ONSFactory.createConsumer(properties);
// Subscribe only to messages where region is hangzhou and price is less than 20.
consumer.subscribe("topicA", MessageSelector.bySql("region IS NOT NULL AND price IS NOT NULL AND region = 'hangzhou' AND price < 20"), new MessageListener() {
    public Action consume(Message message, ConsumeContext context) {
        System.out.printf("Receive New Messages: %s %n", message);
        return Action.CommitMessage;
    }
});

Resultado esperado: A mensagem não é entregue. O consumidor exige price menor que 20, mas a mensagem tem price definido como 50.

Assinar todas as mensagens com SQL

Consumer consumer = ONSFactory.createConsumer(properties);
// Set the SQL expression to TRUE to receive all messages.
consumer.subscribe("topicA", MessageSelector.bySql("TRUE"), new MessageListener() {
    public Action consume(Message message, ConsumeContext context) {
        System.out.printf("Receive New Messages: %s %n", message);
        return Action.CommitMessage;
    }
});

Resultado esperado: Todas as mensagens do tópico são entregues.

Incorreto: atributo ausente e falta de verificação de null

Se um produtor não definir um atributo personalizado em uma mensagem e a expressão de filtro do consumidor referenciar esse atributo sem uma verificação IS NOT NULL, a expressão resultará em NULL. A mensagem não será entregue.

Consumer consumer = ONSFactory.createConsumer(properties);
// The product attribute was not set when the message was sent.
// The filter fails and the message is not delivered.
consumer.subscribe("topicA", MessageSelector.bySql("product = 'MQ'"), new MessageListener() {
    public Action consume(Message message, ConsumeContext context) {
        System.out.printf("Receive New Messages: %s %n", message);
        return Action.CommitMessage;
    }
});

Referências