发布、订阅以及完整的 QoS 机制
TsgcWSPClient_MQTT 通过一个 MQTTVersion 属性同时覆盖两个协议版本。QoS 0、1 和 2 的 PUBACK、PUBREC、PUBREL 和 PUBCOMP 交换以事件形式呈现而非被隐藏,还支持保留消息、遗嘱消息(Last Will and Testament)、通配符订阅和可恢复会话。版本 5 增加了原因码、用户属性、主题别名、共享订阅和 AUTH 往返。
七个客户端组件分别支持 MQTT 3.1.1 和 5.0、AMQP 0.9.1、AMQP 1.0、Apache Kafka 线上协议以及 STOMP 1.0 到 1.2,并为 RabbitMQ 和 ActiveMQ 提供代理方言。每一帧都在库内部用 Object Pascal 解析和写入,因此没有 librdkafka、没有 Paho、没有 Qpid,也没有需要随附发布的 DLL。全部七个组件都运行在同一个 TsgcTCPClient 之上,可以明文打开,也可以包裹在 TLS 中。
sgcMQ 是一个独立软件包,内置其所基于的 sgcWebSockets Core 运行时,每份许可证都附带完整源代码,因此协议实现可以在您自己的调试器中单步跟踪。
TsgcTCPClient 承载全部七个组件,因此 TLS、代理和重连只需配置一次。
把协议组件拖放到窗体上,将一个 TsgcTCPClient 赋给它的 Client 属性,接好事件,然后调用协议自身的动作方法。七个组件的命名保持一致,学会一个,其余的也就掌握了大半。
TsgcWSPClient_MQTT 通过一个 MQTTVersion 属性同时覆盖两个协议版本。QoS 0、1 和 2 的 PUBACK、PUBREC、PUBREL 和 PUBCOMP 交换以事件形式呈现而非被隐藏,还支持保留消息、遗嘱消息(Last Will and Testament)、通配符订阅和可恢复会话。版本 5 增加了原因码、用户属性、主题别名、共享订阅和 AUTH 往返。
TsgcWSPClient_AMQP 使用 AMQP 0.9.1,即 RabbitMQ 赖以构建的协议:信道、交换机和队列声明、绑定、消费者、发布者确认、预取 QoS 以及事务。TsgcWSPClient_AMQP1 使用 AMQP 1.0,即 Azure Service Bus 和 Event Hubs 背后的 OASIS 标准:容器、会话、发送和接收链路、基于信用的流量控制、SASL 以及 Claims-Based Security 辅助方法。线上格式不同,对象模型也不同,因此各有一个组件。
TsgcWSPClient_Kafka 通过 TCP 直接使用 Kafka 线上协议,包括 v2 记录批次格式。可以按选定的 acks 级别和 gzip 压缩生产消息;也可以通过消费组消费,协调器发现、加入、同步、心跳和再均衡都已为您处理好;还可以按偏移量读取分区。
TsgcWSPClient_STOMP 处理 SEND、SUBSCRIBE、UNSUBSCRIBE、ACK、NACK、回执、双向心跳,以及 BEGIN、COMMIT 和 ABORT 事务帧。两个派生组件把代理的目的地约定变成具名方法,一个针对 RabbitMQ,一个针对 ActiveMQ。
全部七个客户端及其下层载体均注册支持 Win32、Win64、Linux64、macOS、iOS 和 Android,覆盖 Delphi 7 到 RAD Studio 13 以及 C++ Builder。这个消息组件包没有任何平台限制。
功能矩阵 →每个组件的参考页面:
创建协议组件,赋值载体,挂接事件,调用动作方法。从明文监听器切换到加密监听器,只是改一个端口号并在载体上设置 TLS := True,其上层的一切都不变。
uses
sgcTCP_Client_WS, sgcWebSocket_Classes, sgcWebSocket_Protocols;
var
TCPClient: TsgcTCPClient;
MQTT: TsgcWSPClient_MQTT;
begin
TCPClient := TsgcTCPClient.Create(nil); // plain MQTT over TCP
TCPClient.Host := 'broker.example.com';
TCPClient.Port := 1883;
TCPClient.WatchDog.Enabled := True;
MQTT := TsgcWSPClient_MQTT.Create(nil);
MQTT.Client := TCPClient;
MQTT.MQTTVersion := mqtt5;
MQTT.Authentication.Enabled := True;
MQTT.Authentication.UserName := 'sgc';
MQTT.Authentication.Password := 'sgc';
MQTT.LastWillTestament.Enabled := True;
MQTT.LastWillTestament.Topic := 'devices/sensor-01/status';
MQTT.LastWillTestament.Message := 'offline';
MQTT.LastWillTestament.QoS := mtqsAtLeastOnce;
MQTT.LastWillTestament.Retain := True;
MQTT.OnMQTTConnect := MQTTConnect;
MQTT.OnMQTTPublish := MQTTPublish;
TCPClient.Active := True;
end;
procedure TForm1.MQTTConnect(Connection: TsgcWSConnection;
const Session: Boolean; const ReasonCode: Integer;
const ReasonName: string;
const ConnectProperties: TsgcWSMQTTCONNACKProperties);
begin
MQTT.Subscribe('sensors/+/temperature', mtqsAtLeastOnce);
end;
procedure TForm1.MQTTPublish(Connection: TsgcWSConnection;
aTopic, aText: string;
PublishProperties: TsgcWSMQTTPublishProperties);
begin
Memo1.Lines.Add(aTopic + ' = ' + aText);
end;
// For an encrypted broker: same code, TLS on the carrier.
// TCPClient.Port := 8883; TCPClient.TLS := True;
uses
sgcTCP_Client_WS, sgcWebSocket_Protocols, sgcAMQP_Classes;
var
TCPClient: TsgcTCPClient;
AMQP: TsgcWSPClient_AMQP;
begin
TCPClient := TsgcTCPClient.Create(nil);
TCPClient.Host := 'broker.example.com';
TCPClient.Port := 5672; // 5671 for the amqps listener
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;
uses
sgcTCP_Client_WS, sgcWebSocket_Protocols, sgcKafka_Classes;
var
TCPClient: TsgcTCPClient;
Kafka: TsgcWSPClient_Kafka;
begin
Kafka := TsgcWSPClient_Kafka.Create(nil);
Kafka.KafkaOptions.ClientId := 'my-delphi-app';
Kafka.KafkaOptions.Producer.Acks := kafkaAcksLeader;
Kafka.KafkaOptions.Producer.Compression := kafkaCompressionGzip;
Kafka.KafkaOptions.Consumer.GroupId := 'my-group';
Kafka.KafkaOptions.Consumer.OffsetReset := kafkaOffsetEarliest;
Kafka.OnKafkaMessage := KafkaMessage;
TCPClient := TsgcTCPClient.Create(nil);
Kafka.Client := TCPClient;
TCPClient.Host := '127.0.0.1';
TCPClient.Port := 9092;
TCPClient.Active := True;
// produce a record to a topic
Kafka.Produce('my-topic', 'Hello Kafka', 'key-1');
// consume: subscribe once, then Poll repeatedly (e.g. from a timer)
Kafka.Subscribe(['my-topic']);
end;
procedure TForm1.KafkaMessage(Sender: TObject;
const Message: TsgcKafkaMessage);
begin
Memo1.Lines.Add(Message.GetKeyString + ' = ' + Message.GetValueString);
end;
// fetch records from a timer, commit when a batch is processed
var
Messages: TsgcKafkaMessages;
begin
Messages := Kafka.Poll(1000);
try
if Messages.Count > 0 then
Kafka.CommitSync;
finally
Messages.Free;
end;
end;
uses
sgcTCP_Client_WS, sgcWebSocket_Classes, sgcWebSocket_Protocols,
sgcWebSocket_Protocol_STOMP_Broker_Client,
sgcWebSocket_Protocol_STOMP_RabbitMQ_Client;
var
TCPClient: TsgcTCPClient;
STOMP: TsgcWSPClient_STOMP_RabbitMQ;
begin
TCPClient := TsgcTCPClient.Create(nil);
TCPClient.Host := 'rabbit.example.com';
TCPClient.Port := 61613; // 61614 for the plugin's TLS listener
STOMP := TsgcWSPClient_STOMP_RabbitMQ.Create(nil);
STOMP.Client := TCPClient;
STOMP.Authentication.Enabled := True;
STOMP.Authentication.UserName := 'guest';
STOMP.Authentication.Password := 'guest';
STOMP.OnRabbitMQConnected := RabbitMQConnected;
STOMP.OnRabbitMQMessage := RabbitMQMessage;
TCPClient.Active := True;
end;
procedure TForm1.RabbitMQConnected(Connection: TsgcWSConnection;
Headers: TsgcWSRabbitMQSTOMPHeadersConnected);
begin
// /queue/orders
STOMP.SubscribeQueue('orders');
STOMP.PublishQueue('orders', '{"orderId":12345}');
end;
procedure TForm1.RabbitMQMessage(Connection: TsgcWSConnection;
MessageText: string; Headers: TsgcWSRabbitMQSTOMPHeadersMessage;
Subscription: TsgcWSBrokerSTOMPSubscriptionItem);
begin
Memo1.Lines.Add(Headers.Destination + ': ' + MessageText);
end;
在 Object Pascal 和 C++ Builder 中是同一种形态,其余三个组件也是同一种形态。完整功能矩阵 →
sgcMQ 需要您准备什么,以及字节如何抵达消息代理。两个答案都很简短,而且都写在这个页面上,而不是藏在小字条款里。
协议 组件 TCP TLS
MQTT 3.1.1 / 5.0 TsgcWSPClient_MQTT 1883 8883
AMQP 0.9.1 TsgcWSPClient_AMQP 5672 5671
AMQP 1.0 TsgcWSPClient_AMQP1 5672 5671
Apache Kafka TsgcWSPClient_Kafka 9092 *
STOMP TsgcWSPClient_STOMP 61613 *
STOMP, RabbitMQ TsgcWSPClient_STOMP_RabbitMQ 61613 61614
STOMP, ActiveMQ TsgcWSPClient_STOMP_ActiveMQ 61613 61612
* 即代理自身加密监听器所配置的端口。
以上是约定俗成的默认值。载体的 Port 属性
取决于您的部署实际监听的端口。
组件所需的一切都内置在 sgcMQ 中:sgcWebSockets Core 运行时、TsgcTCPClient 载体、TLS 层和 JSON 辅助工具。一个安装程序,开箱即含完整源代码。
包中的每个协议都运行在原始套接字上,可明文打开或包裹在 TLS 中。这涵盖了此处列出的原生代理端口,也是 MQTT、AMQP、Kafka 和 STOMP 通常的部署方式。如果要改为通过 WebSocket 运行协议,例如 MQTT over WebSocket 或 Web-STOMP,则需要 WebSocket 客户端,它随 sgcWebSockets 一起提供。选择之前请先确认您的代理暴露的是哪种传输方式。
sgcMQ 面向 Delphi 7 到 RAD Studio 13 以及 C++ Builder,覆盖全部六个平台。此组件包没有 .NET 版本。如果您需要在 .NET 中使用消息传递,请使用 sgcWebSockets。
TLSOptions.IOHandler 可为跨平台构建选择 OpenSSL,或在 Windows 上选择 SChannel;WatchDog 会在链路中断后自动重连;HTTP CONNECT 代理穿越和 IPv6 无需额外配置。配置一次,其上的每个协议都会继承。
sgcMQ 实现的是公开发布的协议,而不是封装某家厂商的 SDK,因此消息代理由您选择:RabbitMQ、Apache Kafka、Eclipse Mosquitto、HiveMQ、EMQX、Apache ActiveMQ、Azure Service Bus 和 Event Hubs,或 AWS IoT Core。更换代理只需改一个主机、一个端口和一组凭据。
sgcMQ 是 eSeGeCe 面向 Delphi、C++ Builder 和 .NET 的九个组件库之一。它们共享同样的约定,提供完整源代码,并可免版税部署。
面向 Delphi 和 C++ Builder 的 MQTT、AMQP 0.9.1 和 1.0、Apache Kafka 与 STOMP 客户端组件。独立产品,随附 sgcWebSockets Core 运行时。
了解更多 →面向 Delphi、C++ Builder、Lazarus 和 .NET 的 WebSocket、HTTP/2、MQTT、AMQP、WebRTC、AI 及 30+ API 集成。WebSocket 载体就在这个库中。
了解更多 →面向 Delphi 和 C++ Builder 的 AI、LLM 和 MCP 组件。一个组件即可对接七家 LLM 提供商,另有 MCP、嵌入和语音。独立产品,随附 sgcWebSockets Core 运行时。
了解更多 →受到全球各地 Delphi、C++ Builder、Lazarus 和 .NET 开发者的信赖。
你们的 sgcWebSockets 库非常实用,而且易于设置。请继续保持!
sgcWebSockets 非常出色,你们的支持也是最棒的!
非常感谢你们的帮助和支持,我很喜欢你们的组件。