AMQP 0.9.1 クライアント
TsgcWSPClient_AMQP は、RabbitMQ が前提としているプロトコルである AMQP 0.9.1 を扱います。モデルをありのままに公開しており、チャネルを開き、エクスチェンジとキューを宣言し、ルーティングキーでバインドしてから、パブリッシュとコンシュームを行います。簡略化したファサードで何かを隠すことはないため、ブローカーは自身のドキュメントどおりに動作します。
TsgcWSPClient_AMQP は、RabbitMQ が前提としているプロトコルである AMQP 0.9.1 を扱います。モデルをありのままに公開しており、チャネルを開き、エクスチェンジとキューを宣言し、ルーティングキーでバインドしてから、パブリッシュとコンシュームを行います。簡略化したファサードで何かを隠すことはないため、ブローカーは自身のドキュメントどおりに動作します。
チャネルには任意の文字列で名前を付けられます。追跡が必要な数値 ID ではなく、'ch1' のような名前がすべての呼び出しで同じチャネルを指します。
TsgcWSPClient_AMQP
AMQP 0-9-1
TCP(5672)または TLS(5671)
Delphi、C++ Builder
トランスポートは素の TCP と TLS です。sgcMQ は素の TCP と TLS で接続します。AMQP over WebSocket が必要な場合は、WebSocket クライアントを提供する sgcWebSockets のパッケージが必要です。
宣言はハンドシェイク完了後の OnAMQPConnect の中で行います。以降、配信は 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}");
ほぼすべてのメソッドには、ブローカーの応答を待って返すブロッキング版の ...Ex が用意されています。順番に処理していくセットアップルーチンではこちらが便利です。
OpenChannel と CloseChannel で、1 本の TCP 接続上に複数の論理的なやり取りを多重化します。EnableChannel と DisableChannel はチャネルのフロー制御を行います。
DeclareExchange(channel, name, type) は、direct、fanout、topic、headers のいずれかの型でエクスチェンジを作成します。DeleteExchange で削除します。
DeclareQueue、BindQueue、UnBindQueue、PurgeQueue、DeleteQueue がキューのライフサイクル全体をカバーします。いずれもブローカーの応答を返す ...Ex 版があります。
Consume(channel, queue) でサブスクリプションを開始し、CancelConsume で終了します。配信は OnAMQPBasicDeliver に届き、単発の取得には OnAMQPBasicGetOk と OnAMQPBasicGetEmpty が使われます。
PublishMessage(channel, exchange, routingKey, body) には、テキスト、ストリーム、ヘッダーとプロパティを備えた完全なメッセージオブジェクト用のオーバーロードがあります。ルーティングできなかったメッセージは OnAMQPBasicReturn に返されます。
AckMessage は配信タグを確認し、RejectMessage は任意の再キュー指定付きで拒否します。Recover は未確認のメッセージをすべて再配信するようブローカーに要求します。
SelectTransaction でチャネルをトランザクションモードにし、CommitTransaction または RollbackTransaction で終了します。各ステップは OnAMQPTransactionOk で確認できます。
AMQPOptions は VirtualHost、MaxChannels、MaxFrameSize、Locale を保持します。いずれも接続開始時のハンドシェイクでブローカーとネゴシエートされます。
SetQoS(channel, prefetchSize, prefetchCount, global) は、ブローカーがコンシューマーに送れる未確認メッセージの数を制限します。処理の遅いハンドラーが大量のメッセージに埋もれることを防げます。
HeartBeat は、プロキシ越しやブローカーのタイムアウトを超えてもアイドル接続を維持します。やり取りのたびに OnAMQPHeartBeat が発生します。