AMQP 0.9.1 Client
TsgcWSPClient_AMQP 使用 AMQP 0.9.1,也就是 RabbitMQ 当初围绕其设计的协议。它按模型本来的样子暴露一切:打开一个信道,声明一个交换机和一个队列,用路由键把它们绑定起来,然后发布和消费。没有任何东西被简化的外观层遮盖,因此消息代理的行为与它自己的文档所述完全一致。
TsgcWSPClient_AMQP 使用 AMQP 0.9.1,也就是 RabbitMQ 当初围绕其设计的协议。它按模型本来的样子暴露一切:打开一个信道,声明一个交换机和一个队列,用路由键把它们绑定起来,然后发布和消费。没有任何东西被简化的外观层遮盖,因此消息代理的行为与它自己的文档所述完全一致。
信道用您自己选择的字符串命名,因此 'ch1' 在每次调用中都指向同一个信道,而不是一个需要您自行跟踪的数字 id。
TsgcWSPClient_AMQP
AMQP 0-9-1
TCP (5672) 或 TLS (5671)
Delphi、C++ Builder
传输:纯 TCP 和 TLS。sgcMQ 通过纯 TCP 和 TLS 连接。如果您需要 AMQP over WebSocket,则需要 sgcWebSockets 包,它提供 WebSocket 客户端。
在握手完成之后,于 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 确认一个投递标签,RejectMessage 拒绝它并可选择重新入队,Recover 则请求消息代理重新投递所有仍未确认的消息。
SelectTransaction 把信道置于事务模式,随后由 CommitTransaction 或 RollbackTransaction 结束它。OnAMQPTransactionOk 确认每一步。
AMQPOptions 承载 VirtualHost、MaxChannels、MaxFrameSize 和 Locale,它们都会在开场握手期间与消息代理协商。
SetQoS(channel, prefetchSize, prefetchCount, global) 限制消息代理可以推送给某个消费者的未确认消息数量,从而不会淹没处理较慢的处理程序。
HeartBeat 让空闲连接得以穿过代理并越过消息代理的超时限制保持存活,每次交互都会触发 OnAMQPHeartBeat。