AMQP 0.9.1-Client-Komponente: sgcMQ | eSeGeCe

AMQP 0.9.1-Client

TsgcWSPClient_AMQP spricht AMQP 0.9.1, das Protokoll, um das herum RabbitMQ entworfen wurde. Es legt das Modell so offen, wie es wirklich ist: Du öffnest einen Kanal, deklarierst einen Exchange und eine Queue, bindest beide mit einem Routing Key und veröffentlichst und konsumierst dann. Nichts versteckt sich hinter einer vereinfachten Fassade, der Broker verhält sich also genau so, wie es seine eigene Dokumentation beschreibt.

TsgcWSPClient_AMQP

Kanäle bekommen einen Namen deiner Wahl, sodass 'ch1' in jedem Aufruf denselben Kanal bezeichnet, statt einer numerischen ID, die du selbst nachhalten müsstest.

Komponentenklasse

TsgcWSPClient_AMQP

Spezifikation

AMQP 0-9-1

Transport

TCP (5672) oder TLS (5671)

Sprachen

Delphi, C++ Builder

Transport: reines TCP und TLS. sgcMQ verbindet sich über reines TCP und TLS. Wenn du AMQP über WebSocket brauchst, benötigst du dafür ein sgcWebSockets Paket, das den WebSocket-Client bereitstellt.

Eine Topologie deklarieren, dann veröffentlichen und konsumieren

Nimm die Deklarationen in OnAMQPConnect vor, sobald der Handshake abgeschlossen ist. Zustellungen treffen danach auf OnAMQPBasicDeliver ein.

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}");

Wichtige Eigenschaften & Methoden

Fast jede Methode hat einen blockierenden ...Ex-Zwilling, der auf die Antwort des Brokers wartet und sie zurückgibt, genau das, was du in einer linearen Einrichtungsroutine brauchst.

Kanäle

OpenChannel und CloseChannel multiplexen mehrere logische Konversationen über eine TCP-Verbindung. EnableChannel und DisableChannel wenden die Flusssteuerung des Kanals an.

Exchanges

DeclareExchange(channel, name, type) erzeugt einen Exchange vom Typ direct, fanout, topic oder headers. DeleteExchange entfernt ihn wieder.

Queues

DeclareQueue, BindQueue, UnBindQueue, PurgeQueue und DeleteQueue decken den gesamten Lebenszyklus einer Queue ab, jeweils mit einer ...Ex-Variante, die die Antwort des Brokers zurückgibt.

Konsumieren

Consume(channel, queue) startet ein Abonnement, CancelConsume beendet es. Zustellungen treffen auf OnAMQPBasicDeliver ein, für das Abholen einzelner Nachrichten gibt es OnAMQPBasicGetOk und OnAMQPBasicGetEmpty.

Veröffentlichen

PublishMessage(channel, exchange, routingKey, body) bietet Überladungen für Text, Streams und ein vollständiges Nachrichtenobjekt mit Headern und Properties. Nicht zustellbare Nachrichten kommen über OnAMQPBasicReturn zurück.

Bestätigung

AckMessage bestätigt ein Delivery-Tag, RejectMessage weist es zurück, optional mit erneutem Einreihen, und Recover bittet den Broker, alles noch Unbestätigte erneut zuzustellen.

Transaktionen

SelectTransaction versetzt einen Kanal in den Transaktionsmodus, danach schließt CommitTransaction oder RollbackTransaction ihn ab. OnAMQPTransactionOk bestätigt jeden Schritt.

Verbindungsfeinabstimmung

AMQPOptions trägt VirtualHost, MaxChannels, MaxFrameSize und Locale, die alle beim einleitenden Handshake mit dem Broker ausgehandelt werden.

Prefetch

SetQoS(channel, prefetchSize, prefetchCount, global) begrenzt, wie viele unbestätigte Nachrichten der Broker an einen Consumer schicken darf, damit ein langsamer Handler nicht überflutet wird.

Erreichbarkeit

HeartBeat hält im Leerlauf liegende Verbindungen durch Proxies hindurch und über Broker-Timeouts hinaus am Leben, wobei bei jedem Austausch OnAMQPHeartBeat ausgelöst wird.

Weiter entdecken

Online-HilfeVollständige API-Referenz und Anwendungsleitfaden.
AMQP 1.0-ClientDas andere AMQP, ein anderes Protokoll mit einer eigenen Komponente.
Kostenlose Testversion herunterladenTeste den Client gegen einen lokalen RabbitMQ-Container.
PreiseSingle-, Team- und Site-Lizenzen mit vollständigem Quellcode.
Bestes Preis-Leistungs-Verhältnis: All-AccessAlle eSeGeCe-Produkte, inklusive Premium-Support, ab €1,059 pro Jahr.
All-Access-Preise ansehen

Bereit loszulegen?

Lade die kostenlose Testversion herunter und veröffentliche aus Delphi oder C++ Builder in deinen ersten Exchange.