Apache Kafka에서 스트리밍 데이터 로드하기

Apache Kafka에서 스트리밍 데이터 로드하기 (Load streaming data from Apache Kafka)

Druid의 Kafka indexing service를 사용해 Kafka 스트림에서 Apache Druid로 데이터를 로드하는 방법을 보여 드릴게요. Koalas to the Max 게임의 샘플 중첩(nested) 클릭스트림 데이터를 Kafka 토픽에 로드한 뒤, 그 데이터를 Druid에 수집하는 과정을 단계별로 진행해요.

출처: 문서

본문

사전 준비 (Prerequisites)

이 튜토리얼의 단계를 따라가기 전에, quickstart에서 자동 단일 머신 구성 (automatic single-machine configuration)를 사용해 Druid를 다운로드하고 로컬 머신에서 실행 중이어야 해요. 어떤 데이터도 로드할 필요는 없어요.

Kafka 다운로드하고 시작하기 (Download and start Kafka)

Apache Kafka는 Druid와 잘 어울리는 고처리량(high-throughput) 메시지 버스예요. 이 튜토리얼에서는 Kafka 2.7.0을 사용해요.

  • Kafka를 다운로드하려면 터미널에서 다음 명령을 실행해 주세요:
    curl -O https://archive.apache.org/dist/kafka/2.7.0/kafka_2.13-2.7.0.tgz
    tar -xzf kafka_2.13-2.7.0.tgz
    cd kafka_2.13-2.7.0
    
  • 이 튜토리얼에 사용하는 머신에서 Kafka를 이미 실행 중이라면, /tmp 의 kafka-logs 디렉터리를 삭제하거나 이름을 바꿔 주세요.

Druid와 Kafka는 모두 서비스를 조정·관리하기 위해 Apache ZooKeeper를 사용해요. Druid가 이미 실행 중이므로, Kafka는 시작할 때 Druid ZooKeeper 인스턴스에 연결돼요.

Druid와 Kafka를 다른 머신에서 실행하는 운영 환경이라면, Kafka broker를 시작하기 전에 Kafka ZooKeeper를 먼저 시작해 주세요.

  • Kafka 루트 디렉터리에서 다음 명령을 실행해 Kafka broker를 시작해 주세요:
    ./bin/kafka-server-start.sh config/server.properties
    
  • 새 터미널 창에서 Kafka 루트 디렉터리로 이동해 kttm 이라는 Kafka 토픽을 만드는 다음 명령을 실행해 주세요:
    ./bin/kafka-topics.sh --create --topic kttm --bootstrap-server localhost:9092
    
    Kafka는 토픽을 성공적으로 추가하면 Created topic kttm 메시지를 반환해요.

Kafka에 데이터 로드하기 (Load data into Kafka)

이 섹션에서는 튜토리얼 디렉터리에 샘플 데이터를 다운로드하고 그 데이터를 Kafka 토픽으로 보내요.

  • Kafka 루트 디렉터리에 샘플 데이터용 디렉터리를 만들어 주세요:
    mkdir sample-data
    
  • 새 디렉터리에 샘플 데이터를 다운로드하고 압축을 풀어 주세요:
    (cd sample-data && curl -O https://static.imply.io/example-data/kttm-nested-v2/kttm-nested-v2-2019-08-25.json.gz)
    
  • Kafka 루트 디렉터리에서 다음 명령을 실행해 샘플 이벤트를 kttm Kafka 토픽에 게시해 주세요:
    export KAFKA_OPTS="-Dfile.encoding=UTF-8"
    gzcat ./sample-data/kttm-nested-v2-2019-08-25.json.gz | ./bin/kafka-console-producer.sh --broker-list localhost:9092 --topic kttm
    

Druid에 데이터 로드하기 (Load data into Druid)

이제 Kafka 토픽에 데이터가 있으니, Druid의 Kafka indexing service를 사용해 그 데이터를 Druid에 수집할 수 있어요.

이를 위해 Druid 콘솔 데이터 로더를 사용하거나 supervisor 스펙을 제출할 수 있어요. 각 방법을 시도하려면 아래 단계를 따라 주세요.

콘솔 데이터 로더로 데이터 로드하기 (Load data with the console data loader)

Druid 콘솔 데이터 로더는 supervisor 스펙의 각 섹션을 구성하는 여러 화면을 보여 준 뒤, Kafka 데이터를 수집하는 수집 태스크를 만들어요.

콘솔 데이터 로더를 사용하려면:

  • localhost:8888 로 이동해서 Load data > Streaming 을 클릭해 주세요.
  • Apache Kafka 를 클릭한 뒤 Connect data 를 클릭해 주세요.
  • 부트스트랩 서버에 localhost:9092, 토픽에 kttm 을 입력한 뒤 Apply 를 클릭하고, 다음과 비슷한 데이터가 보이는지 확인해 주세요.
  • Next: Parse data 를 클릭해 주세요. 데이터 로더가 데이터에 맞는 올바른 입력 형식을 자동으로 결정하려고 해요. 샘플 데이터의 경우 입력 형식 json 을 선택해요. 다양한 옵션을 실험해서 Druid가 데이터를 어떻게 파싱하는지 미리 볼 수 있어요.
  • json 입력 형식이 선택된 상태에서 Next: Parse time 을 클릭해 주세요. 먼저 Apply 를 클릭해야 할 수도 있어요. Druid 아키텍처는 기본 타임스탬프 컬럼을 지정하도록 요구해요. Druid는 타임스탬프를 Druid 데이터소스의 __time 컬럼에 저장해요. 운영 환경에서 데이터에 타임스탬프가 없다면, Parse timestamp from: None 을 선택해 플레이스홀더 값을 사용할 수 있어요. 샘플 데이터의 경우 데이터 로더는 원본 데이터의 타임스탬프 컬럼을 기본 시간 컬럼으로 선택해요.
  • Next: ... 를 세 번 클릭해 Transform과 Filter 단계를 지나 Configure schema 로 이동해 주세요. 변환과 필터를 적용하는 것은 이 튜토리얼의 범위를 벗어나므로 이 두 단계에서는 아무것도 입력할 필요가 없어요.
  • Configure schema 단계에서 컬럼의 데이터 타입을 선택하고 Druid에 수집할 dimensions와 metrics를 구성할 수 있어요. 콘솔이 대부분을 처리해 줘요. event, agent, geo_ip dimension이 json 타입이라는 점에 주목해 주세요.
  • Next: Partition 을 클릭해 Druid가 데이터를 세그먼트로 파티셔닝하는 방법을 구성해 주세요.
  • Segment granularity 로 day 를 선택해 주세요. 데이터셋이 작으므로 추가 조정은 필요 없어요. Next: Tune 을 클릭해 Druid가 데이터를 수집하는 방식을 미세 조정해 주세요.
  • Input tuning 에서 Use earliest offset 을 True 로 설정해 주세요 — 스트림의 처음부터 데이터를 소비하고 싶기 때문에 아주 중요해요. 다른 변경은 없으므로 Next: Publish 를 클릭해 주세요.
  • 데이터소스 이름을 kttm-kafka 로 정하고 Next: Edit spec 을 클릭해 스펙을 검토해 주세요. 콘솔이 구성한 스펙을 보여 줘요. 스펙 위의 버튼을 클릭해 이전 단계를 변경하고 그 변경이 스펙에 어떻게 반영되는지 확인할 수 있어요. 스펙을 직접 수정하고 이전 단계에 반영되는 것을 볼 수도 있어요.
  • Submit 을 클릭해 수집 태스크를 만들어 주세요. Druid가 새로 생성된 supervisor에 초점을 맞춘 태스크 뷰를 표시해요. 태스크 뷰는 자동으로 새로고침되므로, supervisor가 태스크를 시작할 때까지 기다려 주세요. Druid가 데이터 수집을 시작하면서 상태가 Pending 에서 Running 으로 바뀌어요.
  • 헤더에서 Datasources 뷰로 이동해 주세요. kttm-kafka 데이터소스가 여기에 나타나면 쿼리할 수 있어요. 자세한 내용은 Query your data를 참고해 주세요.

1분 후에도 데이터소스가 나타나지 않는다면, supervisor가 스트림의 처음부터 데이터를 읽도록 설정하지 않았을 수 있어요 —바로 Tune 단계의 Use earliest offset 설정이에요. Ingestion 페이지로 가서 Actions(...) 메뉴로 supervisor를 종료해 주세요. 샘플 데이터를 다시 로드하고 Tune 단계에서 올바른 설정을 적용해 주세요.

Supervisor 스펙 제출하기 (Submit a supervisor spec)

데이터 로더를 사용하는 대신, supervisor 스펙을 Druid에 제출할 수도 있어요. 콘솔에서 하거나 Druid API로 할 수 있어요.

콘솔 사용하기 (Use the console)

Druid 콘솔로 supervisor 스펙을 제출하려면:

  • 콘솔에서 Ingestion 을 클릭한 뒤, 새로고침 버튼 옆의 말줄임표를 클릭하고 Submit JSON supervisor 를 선택해 주세요.
  • 이 스펙을 JSON 창에 붙여 넣고 Submit 을 클릭해 주세요:
    {
      "type" : "kafka",
      "spec" : {
        "ioConfig" : {
          "type" : "kafka",
          "consumerProperties" : {
            "bootstrap.servers" : "localhost:9092"
          },
          "topic" : "kttm",
          "inputFormat" : {
            "type" : "json"
          },
          "useEarliestOffset" : true
        },
        "tuningConfig" : {
          "type" : "kafka"
        },
        "dataSchema" : {
          "dataSource" : "kttm-kafka-supervisor-console",
          "timestampSpec" : {
            "column" : "timestamp",
            "format" : "iso"
          },
          "dimensionsSpec" : {
            "dimensions" : [
              "session",
              "number",
              "client_ip",
              "language",
              "adblock_list",
              "app_version",
              "path",
              "loaded_image",
              "referrer",
              "referrer_host",
              "server_ip",
              "screen",
              "window",
              { "type" : "long", "name" : "session_length" },
              "timezone",
              "timezone_offset",
              { "type" : "json", "name" : "event" },
              { "type" : "json", "name" : "agent" },
              { "type" : "json", "name" : "geo_ip" }
            ]
          },
          "granularitySpec" : {
            "queryGranularity" : "none",
            "rollup" : false,
            "segmentGranularity" : "day"
          }
        }
      }
    }
    
    이 스펙은 supervisor를 시작해요 — supervisor는 들어오는 데이터를 듣기 시작하는 태스크를 생성해요.
  • 콘솔 홈 페이지에서 Tasks 를 클릭해 작업 상태를 모니터링해 주세요. 이 스펙은 kttm 토픽의 데이터를 kttm-kafka-supervisor-console 이라는 데이터소스에 써요.
API 사용하기 (Use the API)

Druid API로 supervisor 스펙을 제출할 수도 있어요.

  • 다음 명령을 실행해 샘플 스펙을 다운로드해 주세요:
    curl -o kttm-kafka-supervisor.json https://raw.githubusercontent.com/apache/druid/master/docs/assets/files/kttm-kafka-supervisor.json
    
  • kttm-kafka-supervisor.json 파일의 스펙을 제출하려면 다음 명령을 실행해 주세요:
    curl -X POST -H 'Content-Type: application/json' -d @kttm-kafka-supervisor.json http://localhost:8081/druid/indexer/v1/supervisor
    
    Druid가 supervisor를 성공적으로 만들면 supervisor ID를 담은 응답을 받아요: {"id":"kttm-kafka-supervisor-api"}.
  • 콘솔 홈 페이지에서 Tasks 를 클릭해 작업 상태를 모니터링해 주세요. 이 스펙은 kttm 토픽의 데이터를 kttm-kafka-supervisor-api 라는 데이터소스에 써요.

데이터 쿼리하기 (Query your data)

Druid가 Kafka 스트림에서 데이터를 보낸 뒤에는 즉시 쿼리할 수 있어요. Druid 콘솔에서 Query 를 클릭해 데이터소스에 대해 SQL 쿼리를 실행해 주세요.

이 튜토리얼은 작은 데이터셋을 수집하므로, SELECT * FROM "kttm-kafka" 쿼리를 실행해 만든 데이터셋의 모든 데이터를 반환할 수 있어요.

새로 로드한 데이터에 몇 가지 예제 쿼리를 실행하려면 Querying data 튜토리얼을 확인해 주세요.

더 읽어보기 (Further reading)

자세한 내용은 다음 주제를 참고해 주세요:

  • Apache Kafka ingestion — Kafka 스트림에서 데이터를 로드하고 Druid용 Kafka supervisor를 유지 관리하는 방법에 대한 정보.

더 알아보기 (Learn more)