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);
Todas as instâncias de consumidor no mesmo ID de grupo devem usar assinaturas consistentes.
O modo broadcasting não permite envio ou recebimento de mensagens ordenadas, rastreamento do progresso de consumo nem redefinição de offsets do consumidor. Para mais informações, consulte Consumo por clustering e consumo por broadcasting.
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 |
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 |
|
|
ID do grupo de consumidores |
Console do ApsaraMQ for RocketMQ |
|
|
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)
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 |
|
|
Quantidade máxima de mensagens por lote. O SDK invoca o callback quando a contagem em cache atinge esse valor. |
1 a 1024 |
32 |
|
|
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.