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"