첫 스트림 수집

첫 스트림 수집 (First Stream Ingest)

Kafka에서 실시간 스트리밍 수집을 설정하고 데이터가 Pinot에 도착하는 것을 지켜보는 페이지예요. 메시지가 Kafka 토픽에 도착하면 몇 초 안에 Pinot가 읽어서 쿼리 가능한 행으로 만들어 주는 흐름을 직접 확인해 볼 수 있어요.

출처: First Stream Ingest

본문

💡 Kubernetes 전용 스트리밍 수집은 Stream ingestion (Kubernetes)을 참고하세요.

이 페이지를 마치면 알게 되는 것 (Outcome)

이 페이지를 끝까지 읽으면 Kafka 토픽에서 데이터를 소비하는 실시간 테이블이 생기고, 쿼리 콘솔에서 12개 행을 볼 수 있어요.

사전 요구 사항 (Prerequisites)

단계 (Steps)

1. 스트리밍 수집 이해하기

스트리밍 수집은 Pinot가 메시지 큐에서 데이터를 실시간으로 소비할 수 있게 해줘요. Kafka 토픽에 메시지가 도착하면 Pinot가 읽고 몇 초 안에 행을 쿼리 가능하게 만들어요. 실시간 테이블 설정(config)은 Kafka 브로커, 토픽, 디코더를 지정해서 Pinot가 들어오는 레코드에 어떻게 연결하고 해석할지 알게 해줘요.

2. Kafka 시작

로컬:

Pinot 퀵스타트의 같은 ZooKeeper를 사용해 9876 포트에서 Kafka를 시작하세요:

bin/pinot-admin.sh StartKafka -zkAddress=localhost:2123/kafka -port 9876

Docker:

Kafka 4.0은 KRaft 모드로 실행되며 ZooKeeper가 필요 없어요:

docker run \
    --network pinot-demo --name=kafka \
    -e KAFKA_NODE_ID=1 \
    -e KAFKA_PROCESS_ROLES=broker,controller \
    -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 \
    -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092 \
    -e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \
    -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT \
    -e KAFKA_CONTROLLER_QUORUM_VOTERS=1@kafka:9093 \
    -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
    -e CLUSTER_ID=MkU3OEVBNTcwNTJENDM2Qk \
    -d apache/kafka:4.0.0

3. Kafka 토픽 생성

로컬:

아직 Apache Kafka를 다운로드하지 않았다면 받고, 토픽을 만드세요:

bin/kafka-topics.sh --create --bootstrap-server localhost:9876 \
    --replication-factor 1 --partitions 1 --topic transcript-topic

Docker:

docker exec \
  -t kafka \
  /opt/kafka/bin/kafka-topics.sh \
  --bootstrap-server kafka:9092 \
  --partitions=1 --replication-factor=1 \
  --create --topic transcript-topic

4. 실시간 테이블 설정 저장

/tmp/pinot-quick-start/transcript-table-realtime.json 파일을 만드세요:

로컬:

{
  "tableName": "transcript",
  "tableType": "REALTIME",
  "segmentsConfig": {
    "timeColumnName": "timestampInEpoch",
    "timeType": "MILLISECONDS",
    "schemaName": "transcript",
    "replicasPerPartition": "1"
  },
  "tenants": {},
  "tableIndexConfig": {
    "loadMode": "MMAP",
    "streamConfigs": {
      "streamType": "kafka",
      "stream.kafka.topic.name": "transcript-topic",
      "stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
      "stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
      "stream.kafka.broker.list": "localhost:9876",
      "realtime.segment.flush.threshold.rows": "0",
      "realtime.segment.flush.threshold.time": "24h",
      "realtime.segment.flush.threshold.segment.size": "50M",
      "stream.kafka.consumer.prop.auto.offset.reset": "smallest"
    }
  },
  "metadata": { "customConfigs": {} }
}

Docker:

{
  "tableName": "transcript",
  "tableType": "REALTIME",
  "segmentsConfig": {
    "timeColumnName": "timestampInEpoch",
    "timeType": "MILLISECONDS",
    "schemaName": "transcript",
    "replicasPerPartition": "1"
  },
  "tenants": {},
  "tableIndexConfig": {
    "loadMode": "MMAP",
    "streamConfigs": {
      "streamType": "kafka",
      "stream.kafka.topic.name": "transcript-topic",
      "stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
      "stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
      "stream.kafka.broker.list": "kafka:9092",
      "realtime.segment.flush.threshold.rows": "0",
      "realtime.segment.flush.threshold.time": "24h",
      "realtime.segment.flush.threshold.segment.size": "50M",
      "stream.kafka.consumer.prop.auto.offset.reset": "smallest"
    }
  },
  "metadata": { "customConfigs": {} }
}

💡 Docker 버전은 Kafka와 Pinot 컨테이너가 모두 같은 pinot-demo Docker 네트워크에 있기 때문에 브로커 주소로 kafka:9092를 사용해요.

5. 실시간 테이블 설정 업로드

실시간 테이블이 생성되는 즉시 Pinot가 Kafka 토픽에서 소비를 시작해요.

로컬:

bin/pinot-admin.sh AddTable \
    -schemaFile /tmp/pinot-quick-start/transcript-schema.json \
    -tableConfigFile /tmp/pinot-quick-start/transcript-table-realtime.json \
    -exec

💡 첫 테이블과 스키마를 진행하면서 transcript 스키마를 이미 업로드했다면 -schemaFile 플래그를 생략할 수 있어요. 포함해도 안전해요. Pinot가 동일한 스키마를 다시 만들지 않고 건너뛰니까요.

Docker:

docker run --rm -ti \
    --network=pinot-demo \
    -v /tmp/pinot-quick-start:/tmp/pinot-quick-start \
    --name pinot-streaming-table-creation \
    apachepinot/pinot:${PINOT_VERSION} AddTable \
    -schemaFile /tmp/pinot-quick-start/transcript-schema.json \
    -tableConfigFile /tmp/pinot-quick-start/transcript-table-realtime.json \
    -controllerHost pinot-controller \
    -controllerPort 9000 \
    -exec

💡 설정 중 다른 이름을 사용했다면 pinot-controller를 실제 Pinot controller 컨테이너 이름으로 바꾸세요.

6. 샘플 스트리밍 데이터 저장

/tmp/pinot-quick-start/rawdata/transcript.json 파일을 만드세요:

{"studentID":205,"firstName":"Natalie","lastName":"Jones","gender":"Female","subject":"Maths","score":3.8,"timestampInEpoch":1571900400000}
{"studentID":205,"firstName":"Natalie","lastName":"Jones","gender":"Female","subject":"History","score":3.5,"timestampInEpoch":1571900400000}
{"studentID":207,"firstName":"Bob","lastName":"Lewis","gender":"Male","subject":"Maths","score":3.2,"timestampInEpoch":1571900400000}
{"studentID":207,"firstName":"Bob","lastName":"Lewis","gender":"Male","subject":"Chemistry","score":3.6,"timestampInEpoch":1572418800000}
{"studentID":209,"firstName":"Jane","lastName":"Doe","gender":"Female","subject":"Geography","score":3.8,"timestampInEpoch":1572505200000}
{"studentID":209,"firstName":"Jane","lastName":"Doe","gender":"Female","subject":"English","score":3.5,"timestampInEpoch":1572505200000}
{"studentID":209,"firstName":"Jane","lastName":"Doe","gender":"Female","subject":"Maths","score":3.2,"timestampInEpoch":1572678000000}
{"studentID":209,"firstName":"Jane","lastName":"Doe","gender":"Female","subject":"Physics","score":3.6,"timestampInEpoch":1572678000000}
{"studentID":211,"firstName":"John","lastName":"Doe","gender":"Male","subject":"Maths","score":3.8,"timestampInEpoch":1572678000000}
{"studentID":211,"firstName":"John","lastName":"Doe","gender":"Male","subject":"English","score":3.5,"timestampInEpoch":1572678000000}
{"studentID":211,"firstName":"John","lastName":"Doe","gender":"Male","subject":"History","score":3.2,"timestampInEpoch":1572854400000}
{"studentID":212,"firstName":"Nick","lastName":"Young","gender":"Male","subject":"History","score":3.6,"timestampInEpoch":1572854400000}

7. Kafka 토픽에 데이터 푸시

로컬:

bin/kafka-console-producer.sh \
    --bootstrap-server localhost:9876 \
    --topic transcript-topic < /tmp/pinot-quick-start/rawdata/transcript.json

Docker:

docker exec -t kafka /opt/kafka/bin/kafka-console-producer.sh \
    --bootstrap-server localhost:9092 \
    --topic transcript-topic < /tmp/pinot-quick-start/rawdata/transcript.json

확인 (Verify)

  1. 브라우저에서 Query Console을 엽니다.
  2. 다음 쿼리를 실행합니다:
SELECT * FROM transcript
  1. 12개 행의 스트리밍 데이터를 볼 수 있어요. Pinot는 Kafka에서 실시간으로 수집하므로, 토픽에 푸시된 지 몇 초 안에 행이 나타나요.

다음 단계 (Next step)

Pinot 테이블에 대해 분석 쿼리를 작성하는 방법을 배우려면 첫 쿼리로 계속 진행하세요.

더 알아보기 (Learn more)