Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Subscribe to messages

Última atualização: Jun 27, 2026

Assine mensagens com o SDK do cliente TCP para Java do ApsaraMQ for RocketMQ. Este tópico aborda modos de assinatura, modos de entrega e exemplos completos de código.

Modos de assinatura

O ApsaraMQ for RocketMQ oferece dois modos de assinatura: clustering e broadcasting.

Clustering (padrão)

No modo clustering, todos os consumidores do mesmo grupo compartilham a carga de mensagens igualmente. Cada mensagem é entregue a apenas um consumidor.

Por exemplo, se um tópico contiver 9 mensagens e o grupo tiver 3 consumidores, cada um receberá 3 mensagens.

// Clustering is the default mode. To set it explicitly:
properties.put(PropertyKeyConst.MessageModel, PropertyValueConst.CLUSTERING);

Broadcasting

No modo broadcasting, cada consumidor do grupo recebe uma cópia completa de todas as mensagens.

Por exemplo, se um tópico contiver 9 mensagens e o grupo tiver 3 consumidores, cada um receberá todas as 9 mensagens.

// Set the subscription mode to broadcasting.
properties.put(PropertyKeyConst.MessageModel, PropertyValueConst.BROADCASTING);
Nota

Modos de entrega

O ApsaraMQ for RocketMQ disponibiliza dois modos de entrega: push e pull.

Modo

Descrição

Mais indicado para

Push

O ApsaraMQ for RocketMQ envia mensagens aos consumidores. Também oferece suporte ao consumo em lote.

A maioria dos casos de uso

Pull

Os consumidores buscam mensagens no ApsaraMQ for RocketMQ, o que proporciona mais opções de recebimento e maior liberdade na obtenção das mensagens.

Cenários que exigem controle refinado sobre a busca de mensagens

Importante
  • Consumidores pull exigem uma instância da Enterprise Platinum Edition.

  • Consumidores pull conectam-se exclusivamente por meio de uma Virtual Private Cloud (VPC).

Para obter a lista completa de métodos e parâmetros do consumidor pull, consulte Métodos e parâmetros.

Código de exemplo

Os exemplos a seguir demonstram consumo via push, push em lote e pull. Para ver mais exemplos, acesse o repositório de códigos do ApsaraMQ for RocketMQ no GitHub.

Substitua os placeholders abaixo pelos valores reais:

Placeholder

Descrição

Onde encontrar

<your-group-id>

ID do grupo de consumidores

Console do ApsaraMQ for RocketMQ

<your-tcp-endpoint>

Endpoint TCP

Seção TCP Endpoint na página Instance Details do console

Todos os exemplos leem as credenciais AccessKey das variáveis de ambiente ALIBABA_CLOUD_ACCESS_KEY_ID e ALIBABA_CLOUD_ACCESS_KEY_SECRET.

Consumidor push

import com.aliyun.openservices.ons.api.Action;
import com.aliyun.openservices.ons.api.ConsumeContext;
import com.aliyun.openservices.ons.api.Consumer;
import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.MessageListener;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;

import java.util.Properties;

public class ConsumerTest {
   public static void main(String[] args) {
       Properties properties = new Properties();
       // The consumer group ID created in the ApsaraMQ for RocketMQ console.
       properties.put(PropertyKeyConst.GROUP_ID, "<your-group-id>");
       // Read AccessKey credentials from environment variables.
       properties.put(PropertyKeyConst.AccessKey, System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"));
       properties.put(PropertyKeyConst.SecretKey, System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
       // The TCP endpoint from the Instance Details page in the console.
       properties.put(PropertyKeyConst.NAMESRV_ADDR, "<your-tcp-endpoint>");
       // Clustering mode (default). Uncomment the line below to switch to broadcasting.
       // properties.put(PropertyKeyConst.MessageModel, PropertyValueConst.CLUSTERING);
       // properties.put(PropertyKeyConst.MessageModel, PropertyValueConst.BROADCASTING);

       Consumer consumer = ONSFactory.createConsumer(properties);
       // Subscribe to specific tags. Use "||" as a separator.
       consumer.subscribe("TopicTestMQ", "TagA||TagB", new MessageListener() {
           public Action consume(Message message, ConsumeContext context) {
               System.out.println("Receive: " + message);
               return Action.CommitMessage;
           }
       });

       // Subscribe to another topic with all tags.
       // To unsubscribe, remove the subscription code and restart the consumer.
       consumer.subscribe("TopicTestMQ-Other", "*", new MessageListener() {
           public Action consume(Message message, ConsumeContext context) {
               System.out.println("Receive: " + message);
               return Action.CommitMessage;
           }
       });

       consumer.start();
       System.out.println("Consumer Started");
   }
}

Consumidor push (consumo em lote)

Importante

O consumo em lote requer a versão 1.8.7.3 ou posterior do SDK do cliente TCP para Java. Para detalhes sobre versões, consulte as Notas de lançamento.

Duas propriedades controlam o comportamento do lote:

Propriedade

Descrição

Valores válidos

Padrão

ConsumeMessageBatchMaxSize

Quantidade máxima de mensagens por lote. O SDK invoca o callback quando a contagem em cache atinge esse valor.

1 a 1024

32

BatchConsumeMaxAwaitDurationInSeconds

Tempo máximo de espera, em segundos, antes de entregar um lote parcial. O SDK aciona o callback ao término dessa duração, independentemente do tamanho do lote.

0 a 450

0

import com.aliyun.openservices.ons.api.Action;
import com.aliyun.openservices.ons.api.ConsumeContext;
import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.batch.BatchConsumer;
import com.aliyun.openservices.ons.api.batch.BatchMessageListener;
import java.util.List;
import java.util.Properties;

import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.tcp.example.MqConfig;

public class SimpleBatchConsumer {

    public static void main(String[] args) {
        Properties consumerProperties = new Properties();
        // The consumer group ID created in the ApsaraMQ for RocketMQ console.
        consumerProperties.setProperty(PropertyKeyConst.GROUP_ID, MqConfig.GROUP_ID);
        // Read AccessKey credentials from environment variables.
        consumerProperties.put(PropertyKeyConst.AccessKey, System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"));
        consumerProperties.put(PropertyKeyConst.SecretKey, System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
        // The TCP endpoint from the Instance Details page in the console.
        consumerProperties.setProperty(PropertyKeyConst.NAMESRV_ADDR, MqConfig.NAMESRV_ADDR);

        // Set the maximum batch size to 128 messages.
        consumerProperties.setProperty(PropertyKeyConst.ConsumeMessageBatchMaxSize, String.valueOf(128));
        // Set the maximum wait time between batches to 10 seconds.
        consumerProperties.setProperty(PropertyKeyConst.BatchConsumeMaxAwaitDurationInSeconds, String.valueOf(10));

        BatchConsumer batchConsumer = ONSFactory.createBatchConsumer(consumerProperties);
        batchConsumer.subscribe(MqConfig.TOPIC, MqConfig.TAG, new BatchMessageListener() {

             @Override
            public Action consume(final List<Message> messages, ConsumeContext context) {
                System.out.printf("Batch-size: %d\n", messages.size());
                // Process messages in batches.
                return Action.CommitMessage;
            }
        });
        // Start the batch consumer.
        batchConsumer.start();
        System.out.println("Consumer start success.");

        // Keep the process alive to continue consuming.
        try {
            Thread.sleep(200000);
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    }
}

Consumidor pull

import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.PullConsumer;
import com.aliyun.openservices.ons.api.TopicPartition;
import java.util.List;
import java.util.Properties;
import java.util.Set;

public class PullConsumerClient {
    public static void main(String[] args){
        Properties properties = new Properties();
        // The consumer group ID created in the ApsaraMQ for RocketMQ console.
        properties.setProperty(PropertyKeyConst.GROUP_ID, "<your-group-id>");
        // Read AccessKey credentials from environment variables.
        properties.put(PropertyKeyConst.AccessKey, System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"));
        properties.put(PropertyKeyConst.SecretKey, System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
        // The TCP endpoint from the Instance Details page in the console.
        properties.put(PropertyKeyConst.NAMESRV_ADDR, "<your-tcp-endpoint>");
        PullConsumer consumer = ONSFactory.createPullConsumer(properties);
        // Start the pull consumer.
        consumer.start();
        // Get all partitions for the topic.
        Set<TopicPartition> topicPartitions = consumer.topicPartitions("topic-xxx");
        // Assign partitions to this consumer.
        consumer.assign(topicPartitions);

        while (true) {
            // Poll for messages with a 3-second timeout.
            List<Message> messages = consumer.poll(3000);
            System.out.printf("Received message: %s %n", messages);
        }
    }
}

Para mais detalhes sobre partições e offsets, consulte Termos.

Nova tentativa de mensagem

Se o consumo da mensagem falhar ou atingir o tempo limite, o ApsaraMQ for RocketMQ reentrega a mensagem automaticamente. Para saber mais, consulte Nova tentativa de mensagem.