Kafka Client
TsgcWSPClient_Kafka는 9092 포트에서 브로커 자체의 바이너리 프로토콜로 Apache Kafka와 통신합니다. 앞단의 REST 프록시도, 아래의 librdkafka도 없으므로 Kafka 프로듀서나 컨슈머가 처음부터 끝까지 디버깅할 수 있는 자체 완결형 실행 파일 하나가 됩니다.
TsgcWSPClient_Kafka는 9092 포트에서 브로커 자체의 바이너리 프로토콜로 Apache Kafka와 통신합니다. 앞단의 REST 프록시도, 아래의 librdkafka도 없으므로 Kafka 프로듀서나 컨슈머가 처음부터 끝까지 디버깅할 수 있는 자체 완결형 실행 파일 하나가 됩니다.
Kafka는 일반 TCP 프로토콜이므로 캐리어는 TsgcTCPClient입니다. 나머지는 모두 KafkaOptions로 설정합니다.
TsgcWSPClient_Kafka
Apache Kafka 와이어 프로토콜, v2 레코드 배치
TCP (9092), TLS 선택
Delphi, C++ Builder
전송: 일반 TCP와 TLS. sgcMQ는 일반 TCP와 TLS로 연결하며, Kafka 와이어 프로토콜에는 그것으로 충분합니다. 다른 sgcMQ 프로토콜은 MQTT over WebSocket이나 AMQP over WebSocket처럼 WebSocket 전송도 정의하는데, 이 경우에는 WebSocket 클라이언트를 제공하는 sgcWebSockets 패키지가 필요합니다.
GroupId를 지정하면 컴포넌트가 컨슈머 그룹에 참여해 파티션을 할당받고 오프셋을 대신 커밋합니다. 비워 두면 그룹 없이 파티션을 직접 읽습니다.
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;
Kafka.OnKafkaProduce := KafkaProduce;
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, 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;
// include: sgcTCP_Client_WS.hpp, sgcWebSocket_Protocols.hpp,
// sgcKafka_Classes.hpp
TsgcWSPClient_Kafka *Kafka = new TsgcWSPClient_Kafka(this);
Kafka->KafkaOptions->ClientId = "my-cbuilder-app";
Kafka->KafkaOptions->Producer->Acks = kafkaAcksLeader;
Kafka->KafkaOptions->Consumer->GroupId = "my-group";
Kafka->KafkaOptions->Consumer->OffsetReset = kafkaOffsetEarliest;
Kafka->OnKafkaMessage = KafkaMessage;
TsgcTCPClient *TCPClient = new TsgcTCPClient(this);
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(ARRAYOFCONST(("my-topic")));
TsgcKafkaMessages *Messages = Kafka->Poll(1000);
try {
if (Messages->Count > 0)
Kafka->CommitSync();
}
__finally {
delete Messages;
}
가장 자주 사용하게 되는 멤버들입니다.
Produce(topic, value, key, partition)는 레코드 하나를 씁니다. ProduceBytes는 키와 값을 TBytes로 받고, ProduceMessages는 배치 전체를 한 파티션으로 보낸 뒤 브로커 응답을 반환합니다.
KafkaOptions.Producer.Acks로 kafkaAcksNone, kafkaAcksLeader, kafkaAcksAll 중에서 선택합니다. Compression은 gzip을 켜고, TimeoutMs는 브로커의 ack 대기 시간을 제한합니다.
Subscribe([topics])로 관심 토픽을 등록하고 Poll(timeoutMs)로 가져옵니다. Poll은 직접 해제해야 하는 TsgcKafkaMessages 목록을 반환하며, 레코드마다 OnKafkaMessage도 발생시킵니다.
Consumer.GroupId를 설정하면 코디네이터 검색, 조인과 동기화, 파티션 할당, 백그라운드 하트비트가 켜집니다. OnKafkaRebalance가 할당 변경을 알려 줍니다.
CommitSync와 CommitOffset으로 명시적으로 커밋하거나, Consumer.AutoCommit에 간격을 지정합니다. GetEarliestOffset, GetLatestOffset, GetCommittedOffset으로 위치를 다시 읽습니다.
FetchMessages(topic, partition, offset, maxBytes)는 컨슈머 그룹을 완전히 우회하고 특정 파티션을 특정 오프셋부터 읽습니다.
Consumer.MinBytes, MaxBytes, MaxPartitionBytes, MaxWaitMs는 지연 시간과 배치 크기를 맞바꾸고, SessionTimeoutMs와 RebalanceTimeoutMs는 그룹 멤버십을 제어합니다.
CreateTopic은 파티션 수와 복제 계수를 받고, DeleteTopic은 토픽을 삭제하며, GetMetadata, ListGroups, DescribeGroups로 클러스터를 조회합니다.
GetApiVersions는 브로커가 지원하는 프로토콜 API 버전을 물어보므로, 서버 버전을 짐작하는 대신 확실하게 분기할 수 있습니다.