sgcMQ 功能矩阵:MQTT、AMQP、Kafka、STOMP | eSeGeCe

sgcMQ 功能矩阵

sgcMQ 的全部能力,对应到七个客户端组件以及它们所实现的规范。每项能力在 Delphi 和 C++ Builder 中的表现完全一致,每份授权都提供完整源代码。点击组件即可进入它自己的页面,查看用法和示例。

MQTT

3.1.1 和 5.0

AMQP

0.9.1 和 1.0

Kafka

原生线路协议

STOMP

1.0 / 1.1 / 1.2

传输与 TLS

纯 TCP 和 TLS

标准与平台

Delphi 7 至 13,C++ Builder

sgcMQ 是自包含的。它内置了 sgcWebSockets Core 运行时,因此它不是附加组件,也无需购买基础授权。

传输:纯 TCP 和 TLS。sgcMQ 通过纯 TCP 和 TLS 连接。如果您需要在 WebSocket 之上运行协议,例如 MQTT over WebSocket 或 AMQP over WebSocket,则需要 sgcWebSockets 包,它提供 WebSocket 客户端。

下表中的每一项能力都通过原始 TCP 载体 TsgcTCPClient 实现,该载体已包含在 sgcMQ 中。

七个面板组件

四个协议家族,注册在 SGC MQ 面板页上。

组件家族描述
MQTT ClientTsgcWSPClient_MQTTMQTTMQTT 3.1.1 和 5.0 的发布/订阅:QoS 0/1/2、保留消息、Last Will and Testament、会话,以及 MQTT 5 属性。
AMQP 0.9.1 ClientTsgcWSPClient_AMQPAMQP信道、交换机、队列和绑定、消费者、确认、预取 QoS 以及事务。
AMQP 1.0 ClientTsgcWSPClient_AMQP1AMQP会话、发送方和接收方链路、基于信用的流控制、SASL 身份验证,以及 Azure CBS 令牌辅助工具。
Kafka ClientTsgcWSPClient_KafkaKafka支持 acks 和 gzip 的生产者、带再平衡和偏移量提交的消费者组,以及主题管理。
STOMP ClientTsgcWSPClient_STOMPSTOMPSTOMP 1.0/1.1/1.2 帧:SEND、SUBSCRIBE、ACK、NACK、回执、心跳和事务。
STOMP RabbitMQ ClientTsgcWSPClient_STOMP_RabbitMQSTOMP为 RabbitMQ 调优的 STOMP:队列、主题、交换机、外部队列和临时队列辅助方法,可走 TCP 或 TLS。
STOMP ActiveMQ ClientTsgcWSPClient_STOMP_ActiveMQSTOMP为 Apache ActiveMQ 调优的 STOMP,提供队列和主题辅助方法,外加一个 ActiveMQ_Options 属性。

MQTT 3.1.1 和 MQTT 5.0

IoT 领域的发布/订阅协议,完整的服务质量机制是暴露出来的,而不是被隐藏起来。

能力3.1.15.0说明
选择协议版本MQTTVersion 可取 mqtt311mqtt5
QoS 0、1 和 2QoS 2 的交互以 OnMQTTPubRecOnMQTTPubRelOnMQTTPubComp 形式暴露,不会被吞掉。
保留消息retain 标志是 Publish 的第四个参数。
Last Will and TestamentLastWillTestament 携带主题、消息、QoS 和 retain。
通配符订阅+ 匹配一层,# 匹配其余整棵树。
会话OnMQTTConnectSession 参数会报告消息代理是否恢复了某个会话。
发布并等待PublishAndWait 会阻塞直到消息代理确认该数据包。
发布流Publish 提供了 TStream 重载,用于二进制负载。
原因码及其名称MQTT 5 在 CONNACK 和 DISCONNECT 上返回数字原因码及其名称。
CONNECT 属性ConnectProperties 可设置会话过期时间、接收上限、最大数据包大小和主题别名上限。
主题别名用一个短整数在线路上替代冗长的主题名,入站数据包会自动解析还原。
用户属性任意键值对可随 CONNECT、PUBLISH、SUBSCRIBE 和 DISCONNECT 一起传输。
共享订阅订阅 $share/<group>/<topic> 可将投递分散到一组客户端上。
增强身份验证通过 AuthOnMQTTAuth 完成 AUTH 数据包往返,用于质询/响应方案。
保活HeartBeat 驱动 PINGREQ,收到响应时触发 OnMQTTPing

AMQP 0.9.1 和 AMQP 1.0

两个只共用名字、其余毫无关系的协议,因此 sgcMQ 为它们各提供一个独立组件,而不是假装一套 API 能同时适配两者。

能力AMQP 0.9.1AMQP 1.0说明
组件TsgcWSPClient_AMQPTsgcWSPClient_AMQP1线路格式不同,对象模型不同,组件也不同。
多路复用单元信道会话OpenChannel 对应 CreateSession
拓扑声明DeclareExchangeDeclareQueueBindQueueAMQP 1.0 是对消息代理上的节点寻址,而不是声明它们。
发送PublishMessage(channel, exchange, routingKey, body)SendMessage(session, link, text)0.9.1 用路由键,1.0 用目标地址。
接收Consume 之后触发 OnAMQPBasicDeliverCreateReceiverLink 之后触发 OnAMQPMessage一旦建立,两者都是推送式的。
确认AckMessage / RejectMessageOnAMQPMessageSentAck0.9.1 确认投递标签,1.0 结算一次投递。
流控制SetQoS 预取、EnableChannel / DisableChannel基于信用(CreditSizeWindowSize两者都能防止快速生产者压垮慢速消费者。
事务SelectTransactionCommitTransactionRollbackTransaction0.9.1 的 tx 类,按信道暴露。
队列维护PurgeQueueDeleteQueueDeleteExchangeUnBindQueueCloseLinkCloseSession每个方法都有等待消息代理回复的阻塞式 ...Ex 变体。
重新投递Recover / RecoverAsync请求消息代理重新投递未确认的消息。
身份验证OnAMQPAuthenticationOnAMQPChallengeAuthentication.AuthTypeOnAMQPSASLAuthentication1.0 提供 SASL ANONYMOUS、PLAIN 和 EXTERNAL。
云端辅助工具CreateCBSLinkPutCBSTokenCreateAzureCbsSasTokenCreateAzureCbsJWT面向 Azure Service Bus 和 Event Hubs 的 Claims-Based Security。
存活检测HeartBeatPingAMQPOptions.IdleTimeout连接时与消息代理协商确定。
连接调优AMQPOptions.VirtualHostMaxChannelsMaxFrameSizeLocaleAMQPOptions.ContainerIdChannelMaxMaxFrameSizeMaxLinksPerSession在开场握手中发送,并与消息代理协商。
预取SetQoS(channel, prefetchSize, prefetchCount, global)AMQPOptions.CreditSize消息代理最多可以有多少条未确认消息处于传输中。
帧检查OnAMQPBeforeReadFrameOnAMQPBeforeWriteFrame在原始帧被处理或发送之前查看甚至改写它们。

直接对话 Apache Kafka

用 Object Pascal 实现的二进制 Kafka 协议,前面没有 REST 代理,下面也没有 librdkafka。

能力API说明
生产一条记录Produce(topic, value, key, partition)key 和 partition 是可选的。ProduceBytes 对两者都接受 TBytes
批量生产ProduceMessages(topic, partition, messages)返回消息代理对整批记录的生产响应。
投递保证KafkaOptions.Producer.AckskafkaAcksNonekafkaAcksLeaderkafkaAcksAll
压缩KafkaOptions.Producer.CompressionkafkaCompressionNonekafkaCompressionGzip
消费Subscribe([topics]) 之后调用 Poll(timeoutMs)Poll 返回一个 TsgcKafkaMessages 列表,同时为每条记录触发 OnKafkaMessage
消费者组KafkaOptions.Consumer.GroupId协调者发现、加入、同步和心跳都已代为处理。OnKafkaRebalance 报告分配变化。
偏移量策略KafkaOptions.Consumer.OffsetReset没有已提交偏移量时使用 kafkaOffsetEarliestkafkaOffsetLatest
提交偏移量CommitSyncCommitOffset(topic, partition, offset)或者设置 Consumer.AutoCommit 并配合 AutoCommitIntervalMs
直接读取分区FetchMessages(topic, partition, offset, maxBytes)完全绕过消费者组。
偏移量查询GetEarliestOffsetGetLatestOffsetGetCommittedOffset按主题和分区查询。
拉取调优Consumer.MinBytesMaxBytesMaxPartitionBytesMaxWaitMs在延迟和批量大小之间做权衡。
主题管理CreateTopicDeleteTopicGetMetadata创建时可指定分区数和副本因子。
消费者组管理ListGroupsDescribeGroups用您自己的工具检视消费者组。
消息代理能力探测GetApiVersions询问消息代理支持哪些协议 API 版本。
SASLkafkaSaslNonekafkaSaslPlainSASL/PLAIN 用户名和密码身份验证。

STOMP 1.0、1.1 和 1.2

一个通用客户端,外加两个针对特定消息代理的派生组件,把 RabbitMQ 和 ActiveMQ 的目的地约定变成具名方法。

能力STOMPRabbitMQActiveMQ说明
发送帧SendPublishExPublishEx目的地、正文、内容类型以及可选的事务。
订阅Subscribe(id, destination)SubscribeExSubscribeEx派生组件额外提供持久化、独占和确认模式参数。
队列辅助方法SubscribeQueuePublishQueueUnSubscribeQueueSubscribeQueuePublishQueueUnSubscribeQueue/queue/ 目的地前缀的封装。
主题辅助方法SubscribeTopicPublishTopicUnSubscribeTopicSubscribeTopicPublishTopicUnSubscribeTopic/topic/ 目的地前缀的封装。
交换机辅助方法SubscribeExchangePublishExchangeRabbitMQ 的 /exchange/ 目的地,配合路由模式。
外部队列SubscribeQueueOutsidePublishQueueOutsideRabbitMQ 的 /amq/queue/,用于在别处声明的队列。
临时回复队列SubscribeTemporaryQueuePublishTemporaryQueueRabbitMQ 的 /temp-queue/ 请求/回复模式。
确认ACK / NACK确认模式按订阅选择(ackAuto 等)。
事务BeginTransactionCommitTransactionAbortTransaction把多个 SEND 和 ACK 帧组成一个原子单元。
回执OnSTOMPReceiptOnRabbitMQReceiptOnActiveMQReceipt消息代理确认它已处理某个帧。
消息事件OnSTOMPMessageOnRabbitMQMessageOnActiveMQMessage提供目的地、正文和帧头部。
心跳HeartBeatPing在 CONNECT 帧中协商,双向生效。
版本协商Versions.V1_0V1_1V1_2声明您接受的版本,由消息代理选定其一。
虚拟主机Options.VirtualHost在 CONNECT 时作为 host 头部发送。
消息代理扩展ActiveMQ_OptionsActiveMQ 专有头部,以已发布属性的形式暴露。

字节如何抵达消息代理

每个协议组件都通过它的 Client 属性挂到一个 TsgcTCPClient 上。这一个载体让您无需改动协议代码即可使用纯 TCP 和 TLS。

领域细节
纯 TCP把一个 TsgcTCPClient 赋给 Client,协议就会在消息代理的原生端口上直接跑在套接字之上。
WebSocketsgcMQ 中不包含。WebSocket 载体,例如 MQTT over WebSocket 或 Web-STOMP,需要随 sgcWebSockets 提供的 WebSocket 客户端。
TLS载体上的 TLSOptions,其中 TLSOptions.IOHandler 可为跨平台构建选择 OpenSSL,或在 Windows 上选择 SChannel。
客户端证书通过同一个 TLSOptions 实现双向 TLS,证书可来自 PEM 文件、PKCS#12 包或 Windows 证书存储。
代理载体支持 HTTP CONNECT 代理,因此位于企业代理之后的消息代理同样可达。
重连载体上的 WatchDog 会在链路断开后重新连接,间隔和重试次数均可配置。
IPv6由载体支持,因此 IPv6 消息代理地址无需任何额外配置。
线程读取在独立线程上运行。设置载体的通知选项,即可为 VCL 和 FMX 代码把事件编排到主线程上。

规范、编译器和目标平台

公开发布的协议,以及在每个受支持编译器上通用的同一份源代码。

领域细节
MQTTOASIS MQTT 3.1.1 和 MQTT 5.0。
AMQP 0.9.1AMQP 0-9-1 规范:信道、交换机、队列、绑定以及 basic 类。
AMQP 1.0OASIS AMQP 1.0,同时也作为 ISO/IEC 19464 发布。
KafkaTCP 之上的 Apache Kafka 线路协议,包括 v2 记录批格式。
STOMPSTOMP 1.0、1.1 和 1.2。
WebSocket不包含。WebSocket 载体需要 sgcWebSockets,它提供 WebSocket 客户端。
TLS通过 OpenSSL 支持 TLS 1.2 和 TLS 1.3,或使用 Windows SChannel。
编译器Delphi 和 C++ Builder 7 至 13。
平台Win32、Win64、Linux64、macOS、iOS 和 Android。
授权独立产品。已内置 sgcWebSockets Core 运行时,并包含完整源代码。
超值之选:All-AccesseSeGeCe 全部产品,含高级支持,每年 €1,059 起。
查看 All-Access 价格

用 sgcMQ 构建

下载免费试用版,从 Delphi 或 C++ Builder 连接到您的消息代理。