Composant client AMQP 0.9.1 : sgcMQ | eSeGeCe

Client AMQP 0.9.1

TsgcWSPClient_AMQP parle AMQP 0.9.1, le protocole autour duquel RabbitMQ a été conçu. Il expose le modèle tel qu'il est vraiment : tu ouvres un canal, tu déclares un exchange et une file d'attente, tu les lies avec une clé de routage, puis tu publies et tu consommes. Rien n'est caché derrière une façade simplifiée, le broker se comporte donc comme sa propre documentation l'annonce.

TsgcWSPClient_AMQP

Les canaux sont nommés par une chaîne de ton choix, si bien que 'ch1' identifie le même canal dans tous les appels, plutôt qu'un identifiant numérique dont tu devrais garder la trace.

Classe du composant

TsgcWSPClient_AMQP

Spécification

AMQP 0-9-1

Transport

TCP (5672) ou TLS (5671)

Langages

Delphi, C++ Builder

Transport : TCP simple et TLS. sgcMQ se connecte en TCP simple et en TLS. Si tu as besoin d'AMQP sur WebSocket, cela nécessite un package sgcWebSockets, qui fournit le client WebSocket.

Déclarer une topologie, puis publier et consommer

Fais les déclarations dans OnAMQPConnect, une fois la poignée de main terminée. Les livraisons arrivent ensuite sur OnAMQPBasicDeliver.

uses
  sgcTCP_Client_WS, sgcWebSocket_Protocols,
  sgcWebSocket_Protocol_AMQP_Client, sgcAMQP_Classes;

var
  TCPClient: TsgcTCPClient;
  AMQP: TsgcWSPClient_AMQP;
begin
  TCPClient := TsgcTCPClient.Create(nil);
  TCPClient.Host := 'broker.example.com';
  TCPClient.Port := 5672;

  AMQP := TsgcWSPClient_AMQP.Create(nil);
  AMQP.Client := TCPClient;
  AMQP.AMQPOptions.VirtualHost := '/';
  AMQP.HeartBeat.Enabled  := True;
  AMQP.HeartBeat.Interval := 30;

  AMQP.OnAMQPConnect      := AMQPConnect;
  AMQP.OnAMQPBasicDeliver := AMQPBasicDeliver;

  TCPClient.Active := True;
end;

procedure TForm1.AMQPConnect(Sender: TObject);
begin
  AMQP.OpenChannel('ch1');
  AMQP.DeclareExchange('ch1', 'orders', 'direct');
  AMQP.DeclareQueue('ch1', 'orders_in');
  AMQP.BindQueue('ch1', 'orders_in', 'orders', 'create');
  AMQP.Consume('ch1', 'orders_in');

  // Publish to the exchange with a routing key
  AMQP.PublishMessage('ch1', 'orders', 'create', '{"id":42}');
end;
#include "sgcTCP_Client_WS.hpp"
#include "sgcWebSocket_Protocols.hpp"
#include "sgcWebSocket_Protocol_AMQP_Client.hpp"

TsgcTCPClient *TCPClient = new TsgcTCPClient(this);
TCPClient->Host = "broker.example.com";
TCPClient->Port = 5672;

TsgcWSPClient_AMQP *AMQP = new TsgcWSPClient_AMQP(this);
AMQP->Client = TCPClient;
AMQP->OnAMQPConnect = AMQPConnect;
AMQP->OnAMQPBasicDeliver = AMQPBasicDeliver;

TCPClient->Active = true;

// In the OnAMQPConnect handler:
AMQP->OpenChannel("ch1");
AMQP->DeclareExchange("ch1", "orders", "direct");
AMQP->DeclareQueue("ch1", "orders_in");
AMQP->BindQueue("ch1", "orders_in", "orders", "create");
AMQP->Consume("ch1", "orders_in");
AMQP->PublishMessage("ch1", "orders", "create", "{\"id\":42}");

Propriétés et méthodes clés

Presque chaque méthode possède une jumelle bloquante ...Ex qui attend la réponse du broker et la renvoie, ce qui convient à une routine de configuration linéaire.

Canaux

OpenChannel et CloseChannel multiplexent plusieurs conversations logiques sur une seule connexion TCP. EnableChannel et DisableChannel appliquent le contrôle de flux du canal.

Exchanges

DeclareExchange(channel, name, type) crée un exchange de type direct, fanout, topic ou headers. DeleteExchange le supprime.

Files d'attente

DeclareQueue, BindQueue, UnBindQueue, PurgeQueue et DeleteQueue couvrent tout le cycle de vie d'une file, chacune avec une variante ...Ex qui renvoie la réponse du broker.

Consommation

Consume(channel, queue) démarre un abonnement et CancelConsume le termine. Les livraisons arrivent sur OnAMQPBasicDeliver, avec OnAMQPBasicGetOk et OnAMQPBasicGetEmpty pour les récupérations à l'unité.

Publication

PublishMessage(channel, exchange, routingKey, body) possède des surcharges pour le texte, les flux et un objet message complet avec en-têtes et propriétés. Les messages non routables reviennent sur OnAMQPBasicReturn.

Acquittement

AckMessage confirme une étiquette de livraison, RejectMessage la refuse avec remise en file facultative, et Recover demande au broker de relivrer tout ce qui reste non acquitté.

Transactions

SelectTransaction met un canal en mode transactionnel, puis CommitTransaction ou RollbackTransaction le referme. OnAMQPTransactionOk confirme chaque étape.

Réglage de la connexion

AMQPOptions porte VirtualHost, MaxChannels, MaxFrameSize et Locale, tous négociés avec le broker pendant la poignée de main d'ouverture.

Prélecture

SetQoS(channel, prefetchSize, prefetchCount, global) limite le nombre de messages non acquittés que le broker peut pousser vers un consommateur, un gestionnaire lent ne peut donc pas être submergé.

Maintien en vie

HeartBeat garde les connexions inactives vivantes à travers les proxys et au-delà des délais d'attente du broker, avec OnAMQPHeartBeat déclenché à chaque échange.

Pour aller plus loin

Aide en ligneRéférence complète de l'API et guide d'utilisation.
Client AMQP 1.0L'autre AMQP, un protocole différent avec son propre composant.
Télécharger l'essai gratuitFais tourner le client contre un conteneur RabbitMQ local.
TarifsLicences Single, Team et Site avec le code source complet.
Meilleur rapport qualité-prix : All-AccessTous les produits eSeGeCe, Support Premium inclus, à partir de €1,059/an.
Voir les tarifs All-Access

Prêt à te lancer ?

Télécharge l'essai gratuit et publie vers ton premier exchange depuis Delphi ou C++ Builder.