Kafka Sink 커넥터

Kafka Sink 커넥터

Kafka sink 커넥터는 Pulsar 토픽에서 메시지를 가져와 Kafka 토픽에 저장하는 커넥터예요. 이 가이드는 Kafka sink 커넥터를 구성하고 사용하는 방법을 설명해요. Pulsar의 데이터를 Apache Kafka로 전달할 때 사용해요.

참고: 모든 Pulsar 커넥터는 download page에서 내려받을 수 있어요.

출처: 문서

본문

구성 (Configuration)

Kafka sink 커넥터의 구성에는 다음과 같은 파라미터가 있어요.

프로퍼티 (Property)

이름 타입 필수 기본값 설명
bootstrapServers String true " " (빈 문자열) Kafka 클러스터에 초기 연결을 설정하기 위한 호스트/포트 쌍의 쉼표 구분 목록이에요.
securityProtocol String false " " (빈 문자열) Kafka 브로커와 통신하는 데 사용하는 프로토콜이에요.
saslMechanism String false " " (빈 문자열) Kafka 클라이언트 연결에 사용되는 SASL 메커니즘이에요.
saslJaasConfig String false " " (빈 문자열) SASL 연결을 위한 JAAS 로그인 컨텍스트 파라미터로, JAAS 구성 파일에서 사용하는 형식이에요.
sslEnabledProtocols String false " " (빈 문자열) SSL 연결에 활성화되는 프로토콜 목록이에요.
sslEndpointIdentificationAlgorithm String false " " (빈 문자열) 서버 인증서로 서버 호스트네임을 검증하는 엔드포인트 식별 알고리즘이에요.
sslTruststoreLocation String false " " (빈 문자열) 트러스트 스토어(trust store) 파일의 위치예요.
sslTruststorePassword String false " " (빈 문자열) 트러스트 스토어 파일의 비밀번호예요.
acks String true " " (빈 문자열) 요청이 완료되기 전에 프로듀서가 리더가 수신하기를 요구하는 acknowledgment 수예요. 이 값은 전송되는 레코드의 내구성을 제어해요.
batchsize long false 16384L Kafka 프로듀서가 브로커로 보내기 전에 레코드를 함께 배치로 묶으려 시도하는 배치 크기예요.
maxRequestSize long false 1048576L Kafka 요청의 최대 크기(바이트)예요.
topic String true " " (빈 문자열) Pulsar에서 메시지를 받는 Kafka 토픽이에요.
keyDeserializationClass String false org.apache.kafka.common.serialization.StringSerializer Kafka 프로듀서가 키를 직렬화하는 데 사용하는 직렬화 클래스예요.
valueDeserializationClass String false org.apache.kafka.common.serialization.ByteArraySerializer Kafka 프로듀서가 값을 직렬화하는 데 사용하는 직렬화 클래스예요. 이 직렬화기는 KafkaAbstractSink의 특정 구현에 의해 설정돼요.
producerConfigProperties Map false " " (빈 문자열) 프로듀서에 전달되는 프로듀서 구성 프로퍼티예요. 참고: 커넥터 구성 파일에 지정된 다른 프로퍼티가 이 구성보다 우선해요.

예제 (Example)

Kafka sink 커넥터를 사용하기 전에 다음 방법 중 하나로 구성 파일을 만들어야 해요.

  • JSON
{
   "configs": {
      "bootstrapServers": "localhost:6667",
      "topic": "test",
      "acks": "1",
      "batchSize": "16384",
      "maxRequestSize": "1048576",
      "producerConfigProperties": {
         "client.id": "test-pulsar-producer",
         "security.protocol": "SASL_PLAINTEXT",
         "sasl.mechanism": "GSSAPI",
         "sasl.kerberos.service.name": "kafka",
         "acks": "all"
      }
   }
}
  • YAML
configs:
    bootstrapServers: "localhost:6667"
    topic: "test"
    acks: "1"
    batchSize: "16384"
    maxRequestSize: "1048576"
    producerConfigProperties:
        client.id: "test-pulsar-producer"
        security.protocol: "SASL_PLAINTEXT"
        sasl.mechanism: "GSSAPI"
        sasl.kerberos.service.name: "kafka"
        acks: "all"

더 알아보기 (Learn more)