sgcMQ 機能一覧: MQTT、AMQP、Kafka、STOMP | eSeGeCe

sgcMQ 機能一覧

sgcMQ でできることすべてを、7 つのクライアントコンポーネントと、それらが実装する仕様に対応付けて示します。すべての機能は Delphi と C++ Builder で同じように動作し、すべてのライセンスに完全なソースコードが付属します。各コンポーネントをクリックすると、使い方とサンプルを掲載した個別ページに移動します。

MQTT

3.1.1 と 5.0

AMQP

0.9.1 と 1.0

Kafka

ネイティブのワイヤプロトコル

STOMP

1.0 / 1.1 / 1.2

規格とプラットフォーム

Delphi 7 から 13、C++ Builder

sgcMQ は単体で完結しています。sgcWebSockets Core ランタイムを同梱して出荷されるため、アドオンではなく、購入すべき基本ライセンスもありません。

トランスポートは素の TCP と TLS です。sgcMQ は素の TCP と TLS で接続します。MQTT over WebSocket や AMQP over WebSocket のようにプロトコルを WebSocket 上で動かす必要がある場合は、WebSocket クライアントを提供する sgcWebSockets のパッケージが必要です。

以下の表に挙げた機能はすべて、sgcMQ に同梱されている生の TCP キャリア TsgcTCPClient を通じて利用します。

7 つのパレットコンポーネント

4 つのプロトコルファミリーが、SGC MQ パレットページに登録されます。

コンポーネントクラスファミリー説明
MQTT クライアントTsgcWSPClient_MQTTMQTTMQTT 3.1.1 と 5.0 のパブリッシュ/サブスクライブ。QoS 0/1/2、保持メッセージ、Last Will and Testament、セッション、MQTT 5 のプロパティに対応します。
AMQP 0.9.1 クライアントTsgcWSPClient_AMQPAMQPチャネル、エクスチェンジ、キューとバインディング、コンシューマー、確認応答、プリフェッチ QoS、トランザクションに対応します。
AMQP 1.0 クライアントTsgcWSPClient_AMQP1AMQPセッション、送信リンクと受信リンク、クレジットベースのフロー制御、SASL 認証、Azure CBS トークンのヘルパーに対応します。
Kafka クライアントTsgcWSPClient_KafkaKafkaacks と gzip に対応したプロデューサー、リバランスとオフセットコミットを備えたコンシューマーグループ、トピックの管理操作に対応します。
STOMP クライアントTsgcWSPClient_STOMPSTOMPSTOMP 1.0/1.1/1.2 のフレーム。SEND、SUBSCRIBE、ACK、NACK、レシート、ハートビート、トランザクションに対応します。
STOMP RabbitMQ クライアントTsgcWSPClient_STOMP_RabbitMQSTOMPRabbitMQ 向けに調整した STOMP。キュー、トピック、エクスチェンジ、外部キュー、一時キューのヘルパーを TCP または TLS 上で利用できます。
STOMP ActiveMQ クライアントTsgcWSPClient_STOMP_ActiveMQSTOMPApache ActiveMQ 向けに調整した STOMP。キューとトピックのヘルパーに加え、ActiveMQ_Options プロパティを備えます。

MQTT 3.1.1 と MQTT 5.0

IoT のパブリッシュ/サブスクライブプロトコルを、QoS の仕組みを隠さずそのまま公開して提供します。

機能3.1.15.0備考
プロトコルバージョンの選択MQTTVersionmqtt311 または mqtt5 を受け取ります。
QoS 0、1、2QoS 2 のやり取りは OnMQTTPubRecOnMQTTPubRelOnMQTTPubComp として公開され、内部に隠されません。
保持メッセージretain フラグは Publish の第 4 引数です。
Last Will and TestamentLastWillTestament がトピック、メッセージ、QoS、retain を保持します。
ワイルドカードのサブスクリプション1 階層は +、以降のツリー全体は # です。
セッション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

名前以外に共通点のない 2 つのプロトコルです。1 つの API で両方をまかなうふりをせず、sgcMQ はそれぞれに専用のコンポーネントを提供します。

機能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 が利用できます。
クラウド向けヘルパーCreateCBSLinkPutCBSTokenCreateAzureCbsSasTokenCreateAzureCbsJWTAzure Service Bus と Event Hubs 向けの Claims-Based Security です。
死活監視HeartBeatPingAMQPOptions.IdleTimeout接続時にブローカーとネゴシエートされます。
接続のチューニングAMQPOptions.VirtualHostMaxChannelsMaxFrameSizeLocaleAMQPOptions.ContainerIdChannelMaxMaxFrameSizeMaxLinksPerSession接続開始時のハンドシェイクで送信され、ブローカーとネゴシエートされます。
プリフェッチSetQoS(channel, prefetchSize, prefetchCount, global)AMQPOptions.CreditSizeブローカーが未確認のまま送出できるメッセージ数です。
フレームの検査OnAMQPBeforeReadFrameOnAMQPBeforeWriteFrame処理前や送信前に生フレームを確認したり書き換えたりできます。

Apache Kafka と直接通信

Kafka のバイナリプロトコルを Object Pascal で実装しています。前段の REST プロキシも、下層の librdkafka も不要です。

機能API備考
レコードのプロデュースProduce(topic, value, key, partition)キーとパーティションは省略可能です。ProduceBytes は両方を TBytes で受け取ります。
バッチのプロデュースProduceMessages(topic, partition, messages)バッチ全体に対するブローカーの produce レスポンスを返します。
配信保証KafkaOptions.Producer.AckskafkaAcksNonekafkaAcksLeaderkafkaAcksAll から選択します。
圧縮KafkaOptions.Producer.CompressionkafkaCompressionNone または kafkaCompressionGzip です。
コンシュームSubscribe([topics]) のあと Poll(timeoutMs)PollTsgcKafkaMessages のリストを返すとともに、レコードごとに OnKafkaMessage を発生させます。
コンシューマーグループKafkaOptions.Consumer.GroupIdコーディネーターの検出、参加、同期、ハートビートは自動で処理されます。割り当ての変更は OnKafkaRebalance が通知します。
オフセットの方針KafkaOptions.Consumer.OffsetResetコミット済みオフセットがない場合に kafkaOffsetEarliestkafkaOffsetLatest を選びます。
オフセットのコミットCommitSyncCommitOffset(topic, partition, offset)または Consumer.AutoCommitAutoCommitIntervalMs を設定します。
パーティションの直接読み取りFetchMessages(topic, partition, offset, maxBytes)コンシューマーグループを完全にバイパスします。
オフセットの照会GetEarliestOffsetGetLatestOffsetGetCommittedOffsetトピックとパーティション単位で取得します。
フェッチのチューニングConsumer.MinBytesMaxBytesMaxPartitionBytesMaxWaitMsレイテンシとバッチサイズのバランスを調整します。
トピックの管理CreateTopicDeleteTopicGetMetadata作成時にパーティション数とレプリケーションファクターを指定します。
グループの管理ListGroupsDescribeGroups自作のツールからコンシューマーグループを調べられます。
ブローカーの対応状況の確認GetApiVersionsブローカーが対応しているプロトコル API のバージョンを問い合わせます。
SASLkafkaSaslNonekafkaSaslPlainSASL/PLAIN のユーザー名とパスワードによる認証です。

STOMP 1.0、1.1、1.2

汎用クライアント 1 つに加え、RabbitMQ と ActiveMQ のデスティネーション規約を名前付きメソッドに変換する、ブローカー固有の派生コンポーネントが 2 つあります。

機能STOMPRabbitMQActiveMQ備考
フレームの送信SendPublishExPublishExデスティネーション、本文、コンテンツタイプ、任意のトランザクションを指定します。
サブスクライブSubscribe(id, destination)SubscribeExSubscribeEx派生コンポーネントでは durable、exclusive、ack モードの引数が追加されます。
キューのヘルパーSubscribeQueuePublishQueueUnSubscribeQueueSubscribeQueuePublishQueueUnSubscribeQueue/queue/ のデスティネーション接頭辞をラップしています。
トピックのヘルパーSubscribeTopicPublishTopicUnSubscribeTopicSubscribeTopicPublishTopicUnSubscribeTopic/topic/ のデスティネーション接頭辞をラップしています。
エクスチェンジのヘルパーSubscribeExchangePublishExchangeルーティングパターンを伴う RabbitMQ の /exchange/ デスティネーションです。
外部で宣言されたキューSubscribeQueueOutsidePublishQueueOutside他の場所で宣言済みのキュー向けの、RabbitMQ の /amq/queue/ です。
一時的な返信キューSubscribeTemporaryQueuePublishTemporaryQueueRabbitMQ の /temp-queue/ によるリクエスト/レスポンスパターンです。
確認応答ACK / NACKack モードはサブスクリプションごとに選択します(ackAuto など)。
トランザクションBeginTransactionCommitTransactionAbortTransaction複数の SEND や ACK フレームを 1 つのアトミックな単位にまとめます。
レシートOnSTOMPReceiptOnRabbitMQReceiptOnActiveMQReceiptブローカーが処理したフレームを確認します。
メッセージイベントOnSTOMPMessageOnRabbitMQMessageOnActiveMQMessageデスティネーション、本文、フレームのヘッダーを渡します。
ハートビートHeartBeatPingCONNECT フレームで双方向にネゴシエートされます。
バージョンネゴシエーションVersions.V1_0V1_1V1_2受け入れ可能なバージョンを告知すると、ブローカーが 1 つを選びます。
仮想ホストOptions.VirtualHostCONNECT の host ヘッダーとして送信されます。
ブローカーの拡張ActiveMQ_OptionsActiveMQ 固有のヘッダーを published プロパティとして公開しています。

バイト列がブローカーに届くまで

どのプロトコルコンポーネントも、Client プロパティを通じて TsgcTCPClient に接続します。このキャリア 1 つで、プロトコル側のコードを変えずに素の TCP と TLS の両方を利用できます。

項目内容
素の TCPClientTsgcTCPClient を割り当てると、プロトコルはブローカーのネイティブポート上のソケットで直接動作します。
WebSocketsgcMQ には含まれません。MQTT over WebSocket や Web-STOMP などの WebSocket キャリアには、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 キャリアには、WebSocket クライアントを提供する sgcWebSockets が必要です。
TLSOpenSSL による 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 からブローカーに接続しましょう。