AMQP 0.9.1 Client 组件:sgcMQ | eSeGeCe

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 孪生方法,它会等待消息代理的回复并将其返回,这正是线性初始化流程中所需要的。

信道

OpenChannelCloseChannel 在一条 TCP 连接上多路复用多个逻辑会话。EnableChannelDisableChannel 施加信道级流控制。

交换机

DeclareExchange(channel, name, type) 创建类型为 directfanouttopicheaders 的交换机。DeleteExchange 删除它。

队列

DeclareQueueBindQueueUnBindQueuePurgeQueueDeleteQueue 覆盖队列的整个生命周期,每个都有返回消息代理响应的 ...Ex 变体。

消费

Consume(channel, queue) 启动一个订阅,CancelConsume 结束它。投递通过 OnAMQPBasicDeliver 到达,单条消息拉取则对应 OnAMQPBasicGetOkOnAMQPBasicGetEmpty

发布

PublishMessage(channel, exchange, routingKey, body) 提供针对文本、流以及带头部和属性的完整消息对象的重载。无法路由的消息会通过 OnAMQPBasicReturn 退回。

确认

AckMessage 确认一个投递标签,RejectMessage 拒绝它并可选择重新入队,Recover 则请求消息代理重新投递所有仍未确认的消息。

事务

SelectTransaction 把信道置于事务模式,随后由 CommitTransactionRollbackTransaction 结束它。OnAMQPTransactionOk 确认每一步。

连接调优

AMQPOptions 承载 VirtualHostMaxChannelsMaxFrameSizeLocale,它们都会在开场握手期间与消息代理协商。

预取

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 向您的第一个交换机发布消息。