Assine mensagens do ApsaraMQ for RocketMQ com o SDK de cliente TCP para C/C++. Este tópico explica os dois modos de assinatura, apresenta uma implementação completa de consumidor e aborda o gerenciamento do ciclo de vida do consumidor.
Modos de assinatura
O ApsaraMQ for RocketMQ oferece dois modos de assinatura que determinam a distribuição de mensagens entre as instâncias de consumidor no mesmo group ID.
Assinatura por clustering (padrão)
Cada mensagem é entregue a apenas uma instância de consumidor no grupo. Todos os consumidores identificados pelo mesmo group ID consomem quantidades iguais de mensagens.
Por exemplo, se um tópico contiver nove mensagens e o grupo tiver três instâncias de consumidor, cada instância receberá três mensagens.
// Clustering subscription is the default mode. To set it explicitly:
factoryInfo.setFactoryProperty(ONSFactoryProperty::MessageModel, ONSFactoryProperty::CLUSTERING);
Use a assinatura por clustering quando cada mensagem precisar ser processada apenas uma vez, como no processamento de pedidos ou na distribuição de tarefas.
Assinatura por broadcasting
Todas as mensagens são entregues a todas as instâncias de consumidor no grupo. Cada consumidor recebe uma cópia completa de cada mensagem.
Por exemplo, se um tópico contiver nove mensagens e o grupo tiver três instâncias de consumidor, cada instância receberá todas as nove mensagens.
// Set broadcasting subscription mode.
factoryInfo.setFactoryProperty(ONSFactoryProperty::MessageModel, ONSFactoryProperty::BROADCASTING);
Use a assinatura por broadcasting quando todas as instâncias precisarem dos mesmos dados, como na atualização de cache local ou na sincronização de configurações.
Todas as instâncias de consumidor que compartilham o mesmo group ID devem usar assinaturas consistentes. Para mais informações, consulte Consistência de assinatura.
A assinatura por broadcasting não oferece suporte a mensagens ordenadas, manutenção do progresso de consumo nem redefinição de offset do consumidor. Para mais detalhes, consulte Consumo por clustering e consumo por broadcasting.
Código de exemplo
O exemplo a seguir cria um PushConsumer que assina dois tópicos, processa as mensagens recebidas por meio de um callback de listener e gerencia o ciclo de vida do consumidor.
#include "ONSFactory.h"
#include <iostream>
#include <thread>
#include <mutex>
using namespace ons;
// Mutex to synchronize access to shared resources across listener threads.
std::mutex console_mtx;
class ExampleMessageListener : public MessageListener {
public:
Action consume(Message& message, ConsumeContext& context) {
// Return CommitMessage after successful processing.
// Return ReconsumeLater if processing fails or you want to consume
// the message again -- the broker redelivers the message after a
// predefined interval.
std::lock_guard<std::mutex> lk(console_mtx);
std::cout << "Received a message. Topic: " << message.getTopic() << ", MsgId: "
<< message.getMsgID() << std::endl;
return CommitMessage;
}
};
int main(int argc, char* argv[]) {
std::cout << "=======Before consuming messages=======" << std::endl;
ONSFactoryProperty factoryInfo;
// Group ID created in the ApsaraMQ for RocketMQ console.
// ApsaraMQ for RocketMQ instances use group IDs instead of producer IDs
// and consumer IDs. This parameter maintains backward compatibility.
factoryInfo.setFactoryProperty(ONSFactoryProperty::ConsumerId, "GID_XXX");
// Retrieve credentials from environment variables.
// Make sure ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET
// are set before running this program.
factoryInfo.setFactoryProperty(ONSFactoryProperty::AccessKey, getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"));
factoryInfo.setFactoryProperty(ONSFactoryProperty::SecretKey, getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
// Endpoint of the ApsaraMQ for RocketMQ instance.
// Get this value from the ApsaraMQ for RocketMQ console.
factoryInfo.setFactoryProperty(ONSFactoryProperty::NAMESRV_ADDR, "http://xxxxxxxxxxxxxxxx.aliyuncs.com:80");
PushConsumer *consumer = ONSFactory::getInstance()->createPushConsumer(factoryInfo);
// Subscribe to messages with a specific tag in topic-1.
const char* topic_1 = "topic-1";
const char* tag_1 = "tag-1";
// Subscribe to all messages in topic-2 by using the wildcard tag "*".
const char* topic_2 = "topic-2";
const char* tag_2 = "*";
// Register the listener and subscribe to both topics.
ExampleMessageListener * message_listener = new ExampleMessageListener();
consumer->subscribe(topic_1, tag_1, message_listener);
consumer->subscribe(topic_2, tag_2, message_listener);
// Start the consumer. Messages begin arriving after this call.
consumer->start();
// Keep the main thread alive. In production, replace this with
// your application's main loop or signal handler.
std::this_thread::sleep_for(std::chrono::milliseconds(60 * 1000));
consumer->shutdown();
delete message_listener;
std::cout << "=======After consuming messages======" << std::endl;
return 0;
}
Próximos passos
Design de controle de tráfego do cliente RocketMQ — melhores práticas para limitação de taxa de consumidores.