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 토픽에 존재하는 데이터 타입에 따라 토픽에 스키마를 적용해요. keyDeserializationClassvalueDeserializationClass 구성 파라미터에서 데이터 타입을 감지할 수 있어요.

valueDeserializationClassorg.apache.kafka.common.serialization.StringDeserializer라면 Pulsar 토픽에 스키마 타입으로 Schema.STRING()을 설정할 수 있어요.

valueDeserializationClassio.confluent.kafka.serializers.KafkaAvroDeserializer라면 Pulsar는 Confluent Schema Registry®에서 AVRO 스키마를 다운로드해 Pulsar 토픽에 올바르게 설정해요. 이 경우 소스의 consumerConfigProperties 구성 항목 안에 schema.registry.url을 설정해야 해요.

keyDeserializationClassorg.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>

더 알아보기 (Learn more)