1. MQTT
2. AMQP 0.9.1
3. Kafka
4. STOMP
uses
sgcTCP_Client_WS, sgcWebSocket_Classes, sgcWebSocket_Protocols;
var
TCPClient: TsgcTCPClient ;
MQTT: TsgcWSPClient_MQTT ;
begin
TCPClient := TsgcTCPClient .Create(nil );
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 ;
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 ;
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' );
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 ;
Kafka.Produce('my-topic' , 'Hello Kafka' , 'key-1' );
Kafka.Subscribe(['my-topic' ]);
end ;
procedure TForm1.KafkaMessage(Sender: TObject;
const Message: TsgcKafkaMessage );
begin
Memo1.Lines.Add(Message.GetKeyString + ' = ' + Message.GetValueString);
end ;
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 ;
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
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 ;