Kafka Source 커넥터
Kafka Source 커넥터
Kafka source 커넥터는 Kafka 토픽에서 메시지를 가져와 Pulsar 토픽에 저장하는 커넥터예요. 이 가이드는 Kafka source 커넥터를 구성하고 사용하는 방법을 설명해요. Kafka의 데이터를 Pulsar로 가져올 때 사용해요.
참고: 모든 Pulsar 커넥터는 download page에서 내려받을 수 있어요.
출처: 문서
본문
구성 (Configuration)
Kafka source 커넥터의 구성에는 다음과 같은 프로퍼티가 있어요.
프로퍼티 (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 | " " (빈 문자열) | 트러스트 스토어 파일의 비밀번호예요. |
| groupId | String | true | " " (빈 문자열) | 이 컨슈머가 속한 컨슈머 프로세스 그룹을 식별하는 고유 문자열이에요. |
| fetchMinBytes | long | false | 1 | 각 fetch 응답에 기대하는 최소 바이트예요. |
| autoCommitEnabled | boolean | false | true | true로 설정하면 컨슈머의 오프셋이 백그라운드에서 주기적으로 커밋돼요. 이 커밋된 오프셋은 프로세스가 실패했을 때 새 컨슈머가 시작하는 위치로 사용돼요. |
| autoCommitIntervalMs | long | false | 5000 | autoCommitEnabled가 true로 설정된 경우 컨슈머 오프셋이 Kafka에 자동 커밋되는 빈도(밀리초)예요. |
| heartbeatIntervalMs | long | false | 3000 | Kafka의 그룹 관리 기능을 사용할 때 컨슈머에 하트비트를 보내는 간격이에요. 참고: heartbeatIntervalMs는 sessionTimeoutMs보다 작아야 해요. |
| sessionTimeoutMs | long | false | 30000 | Kafka의 그룹 관리 기능을 사용할 때 컨슈머 실패를 감지하는 데 사용하는 타임아웃이에요. |
| topic | String | true | " " (빈 문자열) | Pulsar에 메시지를 보내는 Kafka 토픽이에요. |
| consumerConfigProperties | Map | false | " " (빈 문자열) | 컨슈머에 전달되는 컨슈머 구성 프로퍼티예요. 참고: 커넥터 구성 파일에 지정된 다른 프로퍼티가 이 구성보다 우선해요. |
| keyDeserializationClass | String | false | org.apache.kafka.common.serialization.StringDeserializer | Kafka 컨슈머가 키를 역직렬화하는 데 사용하는 디시리얼라이저 클래스예요. 이 디시리얼라이저는 KafkaAbstractSource의 특정 구현에 의해 설정돼요. |
| valueDeserializationClass | String | false | org.apache.kafka.common.serialization.ByteArrayDeserializer | Kafka 컨슈머가 값을 역직렬화하는 데 사용하는 디시리얼라이저 클래스예요. |
| autoOffsetReset | String | false | earliest | 기본 오프셋 재설정 정책이에요. |
스키마 관리 (Schema Management)
이 Kafka source 커넥터는 Kafka 토픽에 존재하는 데이터 타입에 따라 토픽에 스키마를 적용해요. keyDeserializationClass와 valueDeserializationClass 구성 파라미터에서 데이터 타입을 감지할 수 있어요.
valueDeserializationClass가 org.apache.kafka.common.serialization.StringDeserializer라면 Pulsar 토픽에 스키마 타입으로 Schema.STRING()을 설정할 수 있어요.
valueDeserializationClass가 io.confluent.kafka.serializers.KafkaAvroDeserializer라면 Pulsar는 Confluent Schema Registry®에서 AVRO 스키마를 다운로드해 Pulsar 토픽에 올바르게 설정해요. 이 경우 소스의 consumerConfigProperties 구성 항목 안에 schema.registry.url을 설정해야 해요.
keyDeserializationClass가 org.apache.kafka.common.serialization.StringDeserializer가 아니라면 키가 문자열이 아니라는 뜻이며, Kafka Source는 SEPARATED 인코딩의 KeyValue 스키마 타입을 사용해요. Pulsar는 키에 대해 AVRO 형식을 지원해요.
이 경우 다음 프로퍼티를 가진 Pulsar 토픽을 가질 수 있어요.
- 스키마: SEPARATED 인코딩의 KeyValue 스키마
- 키(Key): Kafka 메시지의 키 내용(base64 인코딩)
- 값(Value): Kafka 메시지의 값 내용
- 키 스키마(KeySchema):
keyDeserializationClass에서 감지한 스키마 - 값 스키마(ValueSchema):
valueDeserializationClass에서 감지한 스키마
토픽 컴팩션과 파티션 라우팅은 Kafka 키를 담고 있는 Pulsar 키를 사용하므로, Kafka에서 가진 것과 같은 값으로 구동돼요.
Pulsar 토픽에서 데이터를 소비할 때 KeyValue 스키마를 사용할 수 있어요. 이렇게 하면 데이터를 올바르게 디코딩할 수 있어요. 원시 키에 접근하려면 Message#getKeyBytes() API를 사용하세요.
예제 (Example)
Kafka source 커넥터를 사용하기 전에 다음 방법 중 하나로 구성 파일을 만들어야 해요.
- JSON
{
"bootstrapServers": "pulsar-kafka:9092",
"groupId": "test-pulsar-io",
"topic": "my-topic",
"sessionTimeoutMs": "10000",
"autoCommitEnabled": false
}
- YAML
configs:
bootstrapServers: "pulsar-kafka:9092"
groupId: "test-pulsar-io"
topic: "my-topic"
sessionTimeoutMs: "10000"
autoCommitEnabled: false
사용법 (Usage)
Kafka source 커넥터를 Pulsar 내장 커넥터로 만들어 standalone 클러스터나 온프레미스(on-premises) 클러스터에서 사용할 수 있어요.
Standalone 클러스터
이 예시는 standalone 모드에서 Kafka source 커넥터를 사용해 Kafka에서 데이터를 가져와 Pulsar 토픽에 쓰는 방법을 설명해요.
전제 조건 (Prerequisites)
- Docker(Community Edition) 설치.
단계 (Steps)
-
Confluent Platform을 다운로드하고 시작해요. 자세한 내용은 로컬 Kafka 서비스 설치 문서를 참고하세요.
-
Pulsar 이미지를 내려받고 standalone 모드로 Pulsar를 시작해요.
docker pull apachepulsar/pulsar:latest
docker run -d -it -p 6650:6650 -p 8080:8080 -v $PWD/data:/pulsar/data --name pulsar-kafka-standalone apachepulsar/pulsar:latest bin/pulsar standalone
- kafka-producer.py 프로듀서 파일을 만들어요.
from kafka import KafkaProducer
producer = KafkaProducer(bootstrap_servers='localhost:9092')
future = producer.send('my-topic', b'hello world')
future.get()
- pulsar-client.py 컨슈머 파일을 만들어요.
import pulsar
client = pulsar.Client('pulsar://localhost:6650')
consumer = client.subscribe('my-topic', subscription_name='my-aa')
while True:
msg = consumer.receive()
print msg
print dir(msg)
print("Received message: '%s'" % msg.data())
consumer.acknowledge(msg)
client.close()
- 다음 파일들을 Pulsar에 복사해요.
docker cp pulsar-io-kafka.nar pulsar-kafka-standalone:/pulsar
docker cp kafkaSourceConfig.yaml pulsar-kafka-standalone:/pulsar/conf
- 새 터미널 창을 열고 Kafka source 커넥터를 로컬 실행 모드로 시작해요.
docker exec -it pulsar-kafka-standalone /bin/bash
./bin/pulsar-admin source localrun \
--archive $PWD/pulsar-io-kafka.nar \
--tenant public \
--namespace default \
--name kafka \
--destination-topic-name my-topic \
--source-config-file $PWD/conf/kafkaSourceConfig.yaml \
--parallelism 1
- 새 터미널 창을 열고 로컬에서 Kafka 프로듀서를 실행해요.
python3 kafka-producer.py
- 새 터미널 창을 열고 로컬에서 Pulsar 컨슈머를 실행해요.
python3 pulsar-client.py
컨슈머 터미널 창에 다음 정보가 나타나요.
Received message: 'hello world'
온프레미스 클러스터 (On-premises cluster)
이 예시는 온프레미스 클러스터에서 Kafka source 커넥터를 만드는 방법을 설명해요.
- Kafka 커넥터의 NAR 패키지를 Pulsar connectors 디렉터리에 복사해요.
cp pulsar-io-kafka-{{connector:version}}.nar $PULSAR_HOME/connectors/pulsar-io-kafka-{{connector:version}}.nar
- 모든 내장 커넥터를 다시 로드해요.
PULSAR_HOME/bin/pulsar-admin sources reload
- Kafka source 커넥터가 목록에 있는지 확인해요.
PULSAR_HOME/bin/pulsar-admin sources available-sources
pulsar-admin sources create명령으로 Pulsar 클러스터에 Kafka source 커넥터를 만들어요.
PULSAR_HOME/bin/pulsar-admin sources create \
--source-config-file <absolute path to kafka-source-config.yaml>