AMQP 0.9.1 Client

TsgcWSPClient_AMQP 使用 AMQP 0.9.1,也就是 RabbitMQ 当初围绕其设计的协议。它按模型本来的样子暴露一切:打开一个信道,声明一个交换机和一个队列,用路由键把它们绑定起来,然后发布和消费。没有任何东西被简化的外观层遮盖,因此消息代理的行为与它自己的文档所述完全一致。

TsgcWSPClient_AMQP

信道用您自己选择的字符串命名,因此 '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。

继续探索

在线帮助完整的 API 参考和使用指南。
AMQP 1.0 Client另一个 AMQP,是一个有自己独立组件的不同协议。
下载免费试用版用该客户端对接本地的 RabbitMQ 容器。
价格Single、Team 和 Site 授权,均含完整源代码。
超值之选:All-AccesseSeGeCe 全部产品,含高级支持,每年 €1,059 起。
查看 All-Access 价格

准备好开始了吗?

下载免费试用版,从 Delphi 或 C++ Builder 向您的第一个交换机发布消息。