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.

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= HangzhouMensagens 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.

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.
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 |
|
|
O atributo não existe. |
|
|
|
O atributo 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. |
|
|
|
Faixa numérica, inclusiva. Equivalente a |
|
|
|
Fora de uma faixa numérica. Equivalente a |
|
|
|
O valor está em um conjunto. Os elementos do conjunto devem ser strings. |
|
|
|
Igual e diferente de. Funciona tanto com números quanto com strings. |
|
|
|
Operadores lógicos. Coloque cada condição entre parênteses. |
|
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
Instâncias de consumidor que usam o mesmo group ID devem assinar os mesmos tópicos. Para mais informações, consulte Consistência de assinatura.
Utilize tópicos e tags para classificar mensagens de diferentes serviços. Para mais informações, consulte Melhores práticas de tópicos e tags.