AMQP 0.9.1 Client
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은 하나의 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는 delivery tag를 확인하고, RejectMessage는 선택적 재큐잉과 함께 거부하며, Recover는 아직 확인되지 않은 모든 메시지를 다시 전달하도록 브로커에 요청합니다.
SelectTransaction은 채널을 트랜잭션 모드로 전환하고, CommitTransaction 또는 RollbackTransaction이 이를 마무리합니다. OnAMQPTransactionOk가 각 단계를 확인해 줍니다.
AMQPOptions는 VirtualHost, MaxChannels, MaxFrameSize, Locale을 담으며, 모두 연결 시작 핸드셰이크에서 브로커와 협상됩니다.
SetQoS(channel, prefetchSize, prefetchCount, global)는 브로커가 컨슈머에 밀어 넣을 수 있는 미확인 메시지 수를 제한하므로, 처리가 느린 핸들러가 밀려드는 메시지에 잠기지 않습니다.
HeartBeat는 프록시를 지나고 브로커 타임아웃을 넘겨서도 유휴 연결을 살아 있게 유지하며, 교환이 일어날 때마다 OnAMQPHeartBeat가 발생합니다.