O ApsaraMQ for RocketMQ oferece dois tipos de consumidor: Push consumer e Simple consumer. Cada tipo gerencia a obtenção de mensagens, a concorrência e as novas tentativas de forma diferente. Escolha o tipo adequado ao seu modelo de processamento e aos requisitos de confiabilidade.
Qual tipo de consumidor usar
|
Cenário |
Tipo recomendado |
Motivo |
|
Tempo de processamento previsível, sem threads personalizadas |
Push consumer |
O SDK gerencia a busca de mensagens, a concorrência e as novas tentativas. Registre um listener, processe cada mensagem e retorne o resultado. |
|
Tempo de processamento variável, fluxos de trabalho personalizados |
Simple consumer |
Sua aplicação controla quando buscar mensagens, como distribuí-las entre threads e quando confirmar a conclusão. |
Alterar o tipo de consumidor não afeta os recursos existentes do ApsaraMQ for RocketMQ nem o processamento de negócios.
Estágios de processamento de mensagens
Ambos os tipos de consumidor seguem um ciclo de vida de três estágios:
Recebimento -- Busca mensagens no servidor.
Processamento -- Executa a lógica de negócios em cada mensagem.
Confirmação -- Reporta o resultado (sucesso ou falha) ao servidor.

Os dois tipos diferem na gestão de cada estágio:
|
Recurso |
Push consumer |
Simple consumer |
|
Interface |
Callback de listener -- implemente a lógica dentro do listener e retorne um resultado |
A aplicação chama operações de API para receber, processar e confirmar mensagens |
|
Concorrência |
Gerenciada pelo SDK |
Gerenciada pela aplicação |
|
Flexibilidade |
Altamente encapsulado, menos flexível |
Operações atômicas, altamente personalizável |
|
Mais indicado para |
Consumo padrão com tempo de processamento previsível |
Fluxos de trabalho personalizados, distribuição assíncrona ou consumo em lote |
|
Classes do SDK |
|
|
Push consumers
Um Push consumer encapsula a lógica de busca de mensagens, gerenciamento de threads e novas tentativas. Registre um listener de mensagens durante a inicialização; o SDK cuidará do restante.
Como funciona
O SDK utiliza internamente um modelo de thread Reactor:
Uma thread interna de long-polling obtém mensagens do servidor de forma assíncrona.
As mensagens são armazenadas em uma fila de cache interna.
O SDK despacha mensagens para as threads de consumo, que invocam seu listener.

Código de exemplo:
// Sample consumption: Use a PushConsumer to consume normal messages.
ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "Your Topic";
FilterExpression filterExpression = new FilterExpression("Your Filter Tag", FilterExpressionType.TAG);
PushConsumer pushConsumer = provider.newPushConsumerBuilder()
// Set the consumer group.
.setConsumerGroup("Your ConsumerGroup")
// Set the endpoint.
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("Your Endpoint").build())
// Set the pre-bound subscription relationship.
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
// Set the message listener.
.setMessageListener(new MessageListener() {
@Override
public ConsumeResult consume(MessageView messageView) {
// Consume the message and return the processing result.
return ConsumeResult.SUCCESS;
}
})
.build();
Resultados do listener
O listener de mensagens deve retornar um dos seguintes resultados:
|
Resultado |
Constante do SDK Java |
Comportamento |
|
Sucesso |
|
O servidor atualiza o progresso do consumo. |
|
Falha |
|
O sistema tenta novamente com base na política de nova tentativa do PushConsumer. |
|
Exceção lançada |
(tratado como falha) |
Mesmo comportamento de nova tentativa de uma falha explícita. |
Comportamento de timeout
Se a lógica de processamento bloquear e impedir a conclusão da mensagem dentro do tempo permitido, o SDK envia forçosamente um resultado de falha e trata a mensagem conforme a política de novas tentativas.
Um timeout faz com que o SDK envie um resultado de falha, mas a thread de processamento atual pode não responder à interrupção e continuar em execução.
Restrições de confiabilidade
O Push consumer determina sucesso ou falha estritamente pelo valor de retorno do listener. Para preservar essa garantia:
Processe sincronamente. Conclua todo o processamento antes de retornar o resultado.
Não redistribua mensagens. Não transfira mensagens para outras threads e retorne um resultado antes que essas threads terminem.
Se o listener retornar sucesso antes da conclusão do processamento e este falhar posteriormente, o servidor considerará a mensagem como consumida e não tentará novamente.
Entrega ordenada de mensagens
Quando um grupo de consumidores usa o modo de consumo ordenado, o Push consumer invoca o listener na ordem estrita das mensagens, sem necessidade de configuração adicional. Para mais informações, consulte Mensagens ordenadas.
A entrega ordenada exige processamento síncrono. A distribuição assíncrona personalizada dentro do listener anula a garantia de ordenação.
Quando usar Push consumers
Os Push consumers são mais adequados quando:
O tempo de processamento é previsível. Durações imprevisíveis acionam timeouts frequentes, causando mensagens duplicadas devido às novas tentativas.
O consumo padrão é suficiente. O SDK controla o modelo de threads e entrega mensagens com throughput máximo. Isso simplifica o desenvolvimento, mas não suporta processamento assíncrono ou controle de taxa personalizado.
Classes do SDK
O ApsaraMQ for RocketMQ fornece duas classes de SDK para Push consumers:
**
PushConsumer** -- Consome mensagens de tópicos padrão (não Lite).**
LitePushConsumer** -- Consome mensagens de tópicos do tipo Lite, com controle de consumo na granularidade do tópico Lite.
Exemplo de PushConsumer
// Consume normal messages with PushConsumer
ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "<your-topic>";
FilterExpression filterExpression = new FilterExpression("<your-filter-tag>", FilterExpressionType.TAG);
PushConsumer pushConsumer = provider.newPushConsumerBuilder()
// Set the consumer group
.setConsumerGroup("<your-consumer-group>")
// Set the access endpoint
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("<your-endpoint>").build())
// Bind the subscription
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
// Register the message listener
.setMessageListener(new MessageListener() {
@Override
public ConsumeResult consume(MessageView messageView) {
// Process the message and return the result
return ConsumeResult.SUCCESS;
}
})
.build();
Substitua os seguintes espaços reservados pelos valores reais:
|
Espaço reservado |
Descrição |
Exemplo |
|
|
Nome do tópico |
order-events |
|
|
Tag de mensagem para filtragem |
payment |
|
|
Nome do grupo de consumidores |
order-service-group |
|
|
Endpoint de acesso ao servidor |
-- |
Exemplo de LitePushConsumer
// Consume normal messages with LitePushConsumer
ClientServiceProvider provider = ClientServiceProvider.loadService();
LitePushConsumer litePushConsumer = provider.newLitePushConsumerBuilder()
// Set the access endpoint
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("<your-endpoint>").build())
// Set the topic
.bindTopic("<your-topic>")
// Set the consumer group
.setConsumerGroup("<your-consumer-group>")
// Register the message listener
.setMessageListener(messageView -> {
// Process the message and return the result
return ConsumeResult.SUCCESS;
})
.build();
// Subscribe to Lite topics
litePushConsumer.subscribeLite("<your-lite-topic-1>");
litePushConsumer.subscribeLite("<your-lite-topic-2>");
Simple consumers
Um Simple consumer fornece operações de API atômicas para processamento de mensagens. Sua aplicação controla diretamente a busca de mensagens, o gerenciamento de threads e a confirmação.
Como funciona
Chame
ReceiveMessagepara obter um lote de mensagens do servidor.Distribua as mensagens para suas threads de negócios para processamento.
Chame
AckMessagepara cada mensagem processada com sucesso.
Se o processamento falhar, não envie uma confirmação. A mensagem fica disponível novamente após a expiração da duração de invisibilidade, acionando uma nova tentativa. Para mais informações, consulte Política de nova tentativa do SimpleConsumer.
Operações de API
|
Operação |
Finalidade |
Parâmetros principais |
|
|
Obter mensagens do servidor |
Tamanho do lote: número de mensagens por solicitação. Duração de invisibilidade da mensagem: tempo máximo de processamento antes da reentrega da mensagem. |
|
|
Confirmar consumo bem-sucedido |
Nenhum |
|
|
Estender o tempo de processamento para uma mensagem já recebida |
Duração de invisibilidade da mensagem: novo valor, geralmente usado quando o processamento leva mais tempo do que o esperado inicialmente. |
O servidor usa armazenamento distribuído; portanto, ReceiveMessage pode retornar um resultado vazio mesmo quando existem mensagens. Para lidar com isso, chame ReceiveMessage novamente ou aumente a concorrência de chamadas.
Tratamento de falhas
A tabela a seguir descreve como diferentes cenários de falha afetam a entrega de mensagens:
|
Cenário de falha |
Comportamento |
|
Falha no processamento (nenhum ACK enviado) |
A mensagem torna-se visível novamente após a expiração da duração de invisibilidade. O servidor a reentrega para nova tentativa. |
|
Processamento excede a duração de invisibilidade |
Igual a nenhum ACK -- a mensagem torna-se visível e é reentregue. Use |
|
O consumidor falha antes de enviar o ACK |
A mensagem é reentregue após a expiração da duração de invisibilidade. |
Entrega ordenada de mensagens
Um Simple consumer processa mensagens ordenadas na ordem de armazenamento. Para um grupo de mensagens que devem permanecer em sequência, não é possível recuperar a próxima mensagem até que a anterior seja processada.
Quando usar Simple consumers
Os Simple consumers são recomendados quando:
O tempo de processamento é imprevisível. Especifique uma duração inicial de invisibilidade da mensagem ao chamar
ReceiveMessagee estenda-a comChangeInvisibleDurationse necessário.Fluxos de trabalho personalizados são necessários. O SDK não impõe nenhum modelo de threads -- implemente distribuição assíncrona, consumo em lote ou qualquer padrão personalizado.
O controle de taxa é importante. Seu código decide quando e com que frequência chamar
ReceiveMessage, fornecendo controle direto sobre o throughput.
Classe do SDK
O ApsaraMQ for RocketMQ fornece uma classe de SDK para Simple consumers: SimpleConsumer. Esta classe não consome mensagens de tópicos Lite.
Exemplo de SimpleConsumer
// Consume normal messages with SimpleConsumer
ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "<your-topic>";
FilterExpression filterExpression = new FilterExpression("<your-filter-tag>", FilterExpressionType.TAG);
SimpleConsumer simpleConsumer = provider.newSimpleConsumerBuilder()
// Set the consumer group
.setConsumerGroup("<your-consumer-group>")
// Set the access endpoint
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("<your-endpoint>").build())
// Bind the subscription
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
.build();
try {
// Pull up to 10 messages, wait up to 30 seconds
List<MessageView> messageViewList = simpleConsumer.receive(10, Duration.ofSeconds(30));
messageViewList.forEach(messageView -> {
System.out.println(messageView);
// Acknowledge each message after successful processing
try {
simpleConsumer.ack(messageView);
} catch (ClientException e) {
e.printStackTrace();
}
});
} catch (ClientException e) {
// Handle failures such as throttling, then retry the receive call
e.printStackTrace();
}
Melhores práticas
Controle o tempo de processamento para Push consumers
Mantenha o processamento de mensagens dentro do limiar de timeout. Timeouts frequentes causam novas tentativas desnecessárias e mensagens duplicadas. Se sua aplicação lida regularmente com tarefas de longa duração, mude para um Simple consumer e defina uma duração apropriada de invisibilidade da mensagem.