Komponent AMQP 0.9.1 Client: sgcMQ | eSeGeCe

AMQP 0.9.1 Client

TsgcWSPClient_AMQP mówi w AMQP 0.9.1, protokole, wokół którego zaprojektowano RabbitMQ. Wystawia model takim, jaki jest naprawdę: otwierasz kanał, deklarujesz wymianę i kolejkę, wiążesz je kluczem routingu, a następnie publikujesz i konsumujesz. Nic nie jest ukryte za uproszczoną fasadą, więc broker zachowuje się dokładnie tak, jak opisuje to jego własna dokumentacja.

TsgcWSPClient_AMQP

Kanały nazywasz dowolnym ciągiem znaków, więc 'ch1' identyfikuje ten sam kanał w każdym wywołaniu, zamiast numerycznego identyfikatora, który musiałbyś śledzić.

Klasa komponentu

TsgcWSPClient_AMQP

Specyfikacja

AMQP 0-9-1

Transport

TCP (5672) lub TLS (5671)

Języki

Delphi, C++ Builder

Transport: zwykły TCP i TLS. sgcMQ łączy się przez zwykły TCP i TLS. Jeśli potrzebujesz AMQP przez WebSocket, wymaga to pakietu sgcWebSockets, który dostarcza klienta WebSocket.

Zadeklaruj topologię, a potem publikuj i konsumuj

Wykonaj deklaracje wewnątrz OnAMQPConnect, po zakończeniu handshake'u. Dostarczone wiadomości pojawią się wtedy w 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}");

Kluczowe właściwości i metody

Niemal każda metoda ma blokującego bliźniaka ...Ex, który czeka na odpowiedź brokera i ją zwraca, a tego właśnie potrzebujesz w liniowej procedurze konfiguracyjnej.

Kanały

OpenChannel i CloseChannel multipleksują kilka logicznych rozmów przez jedno połączenie TCP. EnableChannel i DisableChannel stosują kontrolę przepływu kanału.

Wymiany

DeclareExchange(channel, name, type) tworzy wymianę typu direct, fanout, topic lub headers. DeleteExchange ją usuwa.

Kolejki

DeclareQueue, BindQueue, UnBindQueue, PurgeQueue i DeleteQueue pokrywają cały cykl życia kolejki, każda z wariantem ...Ex, który zwraca odpowiedź brokera.

Konsumowanie

Consume(channel, queue) rozpoczyna subskrypcję, a CancelConsume ją kończy. Dostarczone wiadomości trafiają do OnAMQPBasicDeliver, a OnAMQPBasicGetOk i OnAMQPBasicGetEmpty obsługują pobieranie pojedynczych wiadomości.

Publikowanie

PublishMessage(channel, exchange, routingKey, body) ma przeciążenia dla tekstu, strumieni i pełnego obiektu wiadomości z nagłówkami i właściwościami. Wiadomości bez trasy wracają przez OnAMQPBasicReturn.

Potwierdzanie

AckMessage potwierdza znacznik dostarczenia, RejectMessage odrzuca go z opcjonalnym ponownym kolejkowaniem, a Recover prosi brokera o ponowne dostarczenie wszystkiego, co pozostaje niepotwierdzone.

Transakcje

SelectTransaction przełącza kanał w tryb transakcyjny, a następnie CommitTransaction lub RollbackTransaction go zamyka. OnAMQPTransactionOk potwierdza każdy krok.

Strojenie połączenia

AMQPOptions niesie VirtualHost, MaxChannels, MaxFrameSize i Locale, wszystkie negocjowane z brokerem podczas otwierającego handshake'u.

Prefetch

SetQoS(channel, prefetchSize, prefetchCount, global) ogranicza liczbę niepotwierdzonych wiadomości, jakie broker może wypchnąć do konsumenta, więc wolna obsługa nie zostanie zalana.

Utrzymanie połączenia

HeartBeat utrzymuje bezczynne połączenia przy życiu przez serwery proxy i mimo limitów czasu brokera, a OnAMQPHeartBeat jest zgłaszane przy każdej wymianie.

Poznawaj dalej

Pomoc onlinePełna dokumentacja API i przewodnik użytkowania.
AMQP 1.0 ClientTen drugi AMQP, inny protokół z własnym komponentem.
Pobierz bezpłatną wersję próbnąUruchom klienta wobec lokalnego kontenera RabbitMQ.
CennikLicencje Single, Team i Site z pełnym kodem źródłowym.
Najkorzystniejsza oferta: All-AccessWszystkie produkty eSeGeCe, ze wsparciem Premium w cenie, już od €1,059 rocznie.
Zobacz cennik All-Access

Gotowy, aby zacząć?

Pobierz bezpłatną wersję próbną i opublikuj do swojej pierwszej wymiany z Delphi lub C++ Builder.