Kafka Client 컴포넌트: sgcMQ | eSeGeCe

Kafka Client

TsgcWSPClient_Kafka는 9092 포트에서 브로커 자체의 바이너리 프로토콜로 Apache Kafka와 통신합니다. 앞단의 REST 프록시도, 아래의 librdkafka도 없으므로 Kafka 프로듀서나 컨슈머가 처음부터 끝까지 디버깅할 수 있는 자체 완결형 실행 파일 하나가 됩니다.

TsgcWSPClient_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.AckskafkaAcksNone, kafkaAcksLeader, kafkaAcksAll 중에서 선택합니다. Compression은 gzip을 켜고, TimeoutMs는 브로커의 ack 대기 시간을 제한합니다.

소비

Subscribe([topics])로 관심 토픽을 등록하고 Poll(timeoutMs)로 가져옵니다. Poll은 직접 해제해야 하는 TsgcKafkaMessages 목록을 반환하며, 레코드마다 OnKafkaMessage도 발생시킵니다.

컨슈머 그룹

Consumer.GroupId를 설정하면 코디네이터 검색, 조인과 동기화, 파티션 할당, 백그라운드 하트비트가 켜집니다. OnKafkaRebalance가 할당 변경을 알려 줍니다.

오프셋

CommitSyncCommitOffset으로 명시적으로 커밋하거나, Consumer.AutoCommit에 간격을 지정합니다. GetEarliestOffset, GetLatestOffset, GetCommittedOffset으로 위치를 다시 읽습니다.

파티션 직접 읽기

FetchMessages(topic, partition, offset, maxBytes)는 컨슈머 그룹을 완전히 우회하고 특정 파티션을 특정 오프셋부터 읽습니다.

Fetch 튜닝

Consumer.MinBytes, MaxBytes, MaxPartitionBytes, MaxWaitMs는 지연 시간과 배치 크기를 맞바꾸고, SessionTimeoutMsRebalanceTimeoutMs는 그룹 멤버십을 제어합니다.

관리

CreateTopic은 파티션 수와 복제 계수를 받고, DeleteTopic은 토픽을 삭제하며, GetMetadata, ListGroups, DescribeGroups로 클러스터를 조회합니다.

브로커 기능 확인

GetApiVersions는 브로커가 지원하는 프로토콜 API 버전을 물어보므로, 서버 버전을 짐작하는 대신 확실하게 분기할 수 있습니다.

계속 살펴보기

온라인 도움말전체 API 레퍼런스와 사용 안내서.
sgcMQ 전체 컴포넌트7개 컴포넌트의 전체 기능 매트릭스를 살펴보세요.
무료 체험판 다운로드로컬 Kafka 브로커에 레코드를 생산하고 소비해 보세요.
가격Single, Team, Site 라이선스, 전체 소스 코드 포함.
최고의 가성비: All-Access모든 eSeGeCe 제품과 프리미엄 지원이 포함되어 연 €1,059부터 이용할 수 있어요.
All-Access 가격 보기

시작할 준비가 되셨나요?

무료 체험판을 다운로드하고 Delphi 또는 C++ Builder에서 첫 Kafka 레코드를 생산해 보세요.