Kafka 추출 네임스페이스

Kafka 추출 네임스페이스 (Kafka extraction namespace)

Kafka 토픽을 읽어 차원 값의 이름을 바꿔 주는 kafka-extraction-namespace 확장 기능을 소개할게요. 예를 들어 내부 ID를 사람이 읽기 좋은 형식으로 바꿀 때 유용해요.

출처: 문서

본문

이 Apache Druid 확장 기능을 사용하려면 확장 목록(extensions load list)에 druid-lookups-cached-global과 druid-kafka-extraction-namespace를 포함하세요.

업데이트를 가능한 한 신속하게 반영해야 한다면, 키가 기존 값(old value)이고 메시지가 원하는 새 값(new value)(둘 다 UTF-8)인 Kafka 토픽을 LookupExtractorFactory로 연결할 수 있어요:

{
  "type": "kafka",
  "kafkaTopic": "testTopic",
  "kafkaProperties": {
    "bootstrap.servers": "kafka.service:9092"
  }
}

| 파라미터 | 설명 | 필수 | 기본값 | | kafkaTopic | 데이터를 읽을 Kafka 토픽 | 예 | | | kafkaProperties | Kafka consumer 속성 (bootstrap.servers는 반드시 지정해야 함) | 예 | | | connectTimeout | 초기 연결을 기다리는 시간 | 아니요 | 0 (기다리지 않음) | | isOneToOne | 맵이 일대일인지 여부 (Lookup DimensionSpecs 참고) | 아니요 | false |

kafka-extraction-namespace 확장 기능은 이름/키 쌍을 가진 Apache Kafka 토픽을 읽어 차원 값의 이름을 바꿀 수 있게 해줘요. 대표적인 예로는 ID를 사람이 읽을 수 있는 형식으로 바꾸는 경우가 있어요.

동작 방식 (How it Works)

이 extractor는 구성된 Kafka 토픽을 처음부터 소비하고, 모든 레코드를 내부 맵에 추가해요. Kafka 레코드의 키는 맵의 키로, 레코드의 페이로드는 값으로 사용돼요. 쿼리 시간에는 lookup을 사용해 키를 연결된 값으로 변환할 수 있어요. 쿼리에서 lookup을 구성하고 사용하는 방법은 lookups 를 참고하세요. 키와 값은 모두 lookup extractor에 의해 문자열로 저장돼요.

이 extractor는 토픽에 계속 구독하고 있으므로, 새 레코드가 나타나면 lookup 맵에 추가돼요. 덕분에 lookup 값을 거의 실시간으로 갱신할 수 있어요. 같은 키로 레코드 두 개가 토픽에 추가되면, 더 큰 offset을 가진 레코드가 lookup 맵의 이전 레코드를 대체해요. null 페이로드를 가진 레코드는 tombstone 레코드로 취급되고, 연결된 키는 lookup 맵에서 제거돼요.

이 extractor는 입력 토픽을 KTable 과 매우 비슷하게 취급해요. 따라서 Kafka 토픽을 로그 컴팩션(log compaction) 전략으로 만드는 것이 좋아요. 이렇게 하면 키의 가장 최신 버전이 항상 Kafka에 보존돼요. 보존(retention)과 로그 컴팩션을 제대로 구성하지 않으면 Kafka에서 자동으로 제거된 오래된 키는 더 이상 사용할 수 없게 되고, Druid 서비스가 재시작될 때 손실돼요.

예시 (Example)

country_codes 토픽이 소비되고 있고, 다음 레코드들이 다음 순서로 토픽에 추가된다고 가정해 볼게요:

| Offset | Key | Payload | | 1 | NZ | Nu Zeelund | | 2 | AU | Australia | | 3 | NZ | New Zealand | | 4 | AU | null | | 5 | NZ | Aotearoa | | 6 | CZ | Czechia |

이 입력 토픽을 처음부터 소비하면, 다음 매핑을 담은 lookup 네임스페이스가 만들어져요 (Australia 항목이 추가되었다가 삭제된 것에 주목하세요):

| Key | Value | | NZ | Aotearoa | | CZ | Czechia |

이제 쿼리가 이 extraction namespace를 사용하면, 국가 코드를 쿼리 시간에 전체 국가 이름으로 매핑할 수 있어요.

Tombstone과 레코드 삭제 (Tombstones and Deleting Records)

Kafka lookup extractor는 null Kafka 메시지를 tombstone으로 취급해요. 즉 입력 토픽에서 메시지 페이로드가 null인 레코드는 lookup 맵에서 연결된 키를 제거해서, 사실상 그 키를 삭제해요.

제한사항 (Limitations)

consumer 속성 group.id, auto.offset.reset, enable.auto.commit은 확장 기능이 각각 UUID.randomUUID().toString(), earliest, false로 설정하므로 kafkaProperties에 설정할 수 없어요. 전체 토픽을 Druid 서비스가 처음부터 소비해서 완전한 lookup 값 맵을 만들어야 하기 때문이에요. 이 consumer 속성 중 하나라도 설정하면 extractor가 시작되지 않아요.

현재 Kafka lookup extractor는 전체 Kafka 토픽을 로컬 캐시에 넣어요. on-heap 캐싱을 사용 중이고 Kafka 스트림이 고유 키를 많이 쏟아낸다면, 자바 힙(heap)을 쉽게 가득 채울 수 있어요. Off-heap 캐싱이 이런 우려를 완화해 주지만, 그래도 저장할 수 있는 데이터 양에는 한계가 있어요. 현재는 축출(eviction) 정책이 없어요.

Kafka 이름 변경 기능 테스트 (Testing the Kafka rename functionality)

이 설정을 테스트하려면 다음 producer 콘솔을 통해 Kafka 스트림에 키/값 쌍을 보낼 수 있어요:

./bin/kafka-console-producer.sh --property parse.key=true --property key.separator="->" --broker-list localhost:9092 --topic testTopic

그런 다음 OLD_VAL->NEW_VAL을 새 줄(enter 또는 return)로 구분해 출판하면 이름 변경이 반영돼요.

더 알아보기 (Learn more)