Kafka 커넥터

Kafka 커넥터 (Kafka connector)

이 커넥터는 Apache Kafka 토픽을 Trino의 테이블로 사용할 수 있게 해 줘요. 각 메시지는 Trino에서 행 하나로 표시돼요.

출처: 문서

본문

토픽은 실시간(live)일 수 있어요. 데이터가 도착하면 행이 나타나고, 세그먼트가 드롭되면 사라져요. 단일 쿼리에서 같은 테이블을 여러 번 접근하면(예: 셀프 조인) 이상한 동작이 일어날 수 있어요.

이 커넥터는 워커 전체에서 병렬로 Kafka 토픽의 메시지 데이터를 읽고 써서 상당한 성능 향상을 얻어요. 이 병렬화를 위한 데이터셋 크기는 구성 가능해서 특정 요구에 맞게 조정할 수 있어요.

Kafka 커넥터 튜토리얼도 함께 읽어 보세요.

요구사항 (Requirements)

Kafka에 연결하려면 다음이 필요해요:

  • Kafka 브로커 3.3 이상(KRaft 활성화).
  • Trino 코디네이터와 워커에서 Kafka 노드로 네트워크 접근. 기본 포트는 9092예요.

Confluent 테이블 설명 공급자와 함께 Protobuf 디코더를 사용할 때는 다음 추가 단계를 수행해야 해요:

  • Confluent에서 Confluent 버전 8.1.1용 kafka-protobuf-providerkafka-protobuf-types JAR 파일을 복사해 클러스터의 모든 노드에서 Kafka 커넥터 플러그인 디렉토리(<install directory>/plugin/kafka)에 넣으세요. 플러그인 디렉토리는 설치 방법에 따라 달라요.
  • 그 JAR을 복사해 사용함으로써, 당신은 Confluent가 제공하는 Confluent Community License Agreement의 조건에 동의하게 돼요.

Protobuf와 Confluent 테이블 설명 공급자를 사용하지 않는다면 이 단계는 필요 없어요.

설정 (Configuration)

Kafka 커넥터를 구성하려면 속성을 적절히 바꿔 다음 내용으로 etc/catalog/example.properties 카탈로그 속성 파일을 만드세요.

특수한 인증 방법을 쓸 때처럼 Kafka 클러스터에 접근하기 위해 추가 Kafka 클라이언트 속성을 지정해야 하는 경우가 있어요. 그러려면 kafka.config.resources 속성을 추가해 Kafka 설정 파일을 참조하세요. kafka.properties에 명시적으로 정의하면 설정을 덮어쓸 수 있다는 점에 유의하세요:

connector.name=kafka
kafka.table-names=table1,table2
kafka.nodes=host1:port,host2:port
kafka.config.resources=/etc/kafka-configuration.properties

Kafka 클라이언트 로그 레벨 조정:

org.apache.kafka=WARN

테이블 설명 (Table description)

Kafka 테이블은 JSON 테이블 설명 파일로 정의해요. 일반적인 구조는 다음과 같아요:

{
    "tableName": ...,
    "schemaName": ...,
    "topicName": ...,
    "key": {
        "dataFormat": ...,
        "fields": [
            ...
        ]
    },
    "message": {
        "dataFormat": ...,
        "fields": [
            ...
       ]
    }
}

각 필드는 다음과 같은 속성을 갖습니다:

{
    "name": ...,
    "type": ...,
    "dataFormat": ...,
    "mapping": ...,
    "formatHint": ...,
    "hidden": ...,
    "comment": ...
}
  • name: Trino에서 사용할 컬럼 이름
  • type: Trino 데이터 타입
  • dataFormat: 필드의 바이너리 포맷(LONG, INT, BYTE 등)
  • mapping: 토픽 메시지에서 필드 위치/이름
  • formatHint: custom-date-time 같은 포맷 힌트

raw 포맷

raw 데이터 포맷은 메시지의 바이트 오프셋으로 필드를 매핑해요:

{
  "tableName": "example_table_name",
  "schemaName": "example_schema_name",
  "topicName": "example_topic_name",
  "key": { "..." },
  "message": {
    "dataFormat": "raw",
    "fields": [
      {
        "name": "field1",
        "type": "BIGINT",
        "dataFormat": "LONG",
        "mapping": "0"
      },
      {
        "name": "field2",
        "type": "INTEGER",
        "dataFormat": "INT",
        "mapping": "8"
      },
      {
        "name": "field3",
        "type": "SMALLINT",
        "dataFormat": "LONG",
        "mapping": "12"
      },
      {
        "name": "field4",
        "type": "VARCHAR(6)",
        "dataFormat": "BYTE",
        "mapping": "20:26"
      }
    ]
  }
}

mapping이 바이트 위치를 나타내므로, 삽입은 값 그대로 기록돼요:

INSERT INTO example_raw_table (field1, field2, field3, field4)
  VALUES (123456789, 123456, 1234, 'abcdef');

csv 포맷

csv 포맷은 쉼표로 구분된 메시지를 mapping(0부터 시작하는 컬럼 인덱스)으로 매핑해요:

{
  "tableName": "example_table_name",
  "schemaName": "example_schema_name",
  "topicName": "example_topic_name",
  "key": { "..." },
  "message": {
    "dataFormat": "csv",
    "fields": [
      {
        "name": "field1",
        "type": "BIGINT",
        "mapping": "0"
      },
      {
        "name": "field2",
        "type": "VARCHAR",
        "mapping": "1"
      },
      {
        "name": "field3",
        "type": "BOOLEAN",
        "mapping": "2"
      }
    ]
  }
}
INSERT INTO example_csv_table (field1, field2, field3)
  VALUES (123456789, 'example text', TRUE);

json 포맷

json 포맷은 필드 이름으로 매핑하고, custom-date-time 포맷 힌트로 타임스탬프를 파싱할 수 있어요:

{
  "tableName": "example_table_name",
  "schemaName": "example_schema_name",
  "topicName": "example_topic_name",
  "key": { "..." },
  "message": {
    "dataFormat": "json",
    "fields": [
      {
        "name": "field1",
        "type": "BIGINT",
        "mapping": "field1"
      },
      {
        "name": "field2",
        "type": "VARCHAR",
        "mapping": "field2"
      },
      {
        "name": "field3",
        "type": "TIMESTAMP",
        "dataFormat": "custom-date-time",
        "formatHint": "yyyy-dd-MM HH:mm:ss.SSS",
        "mapping": "field3"
      }
    ]
  }
}
INSERT INTO example_json_table (field1, field2, field3)
  VALUES (123456789, 'example text', TIMESTAMP '2020-07-15 01:02:03.456');

avro 포맷

avro 포맷은 별도의 .avsc 스키마 파일을 참조해요:

{
  "tableName": "example_table_name",
  "schemaName": "example_schema_name",
  "topicName": "example_topic_name",
  "key": { "..." },
  "message":
  {
    "dataFormat": "avro",
    "dataSchema": "/avro_message_schema.avsc",
    "fields":
    [
      {
        "name": "field1",
        "type": "BIGINT",
        "mapping": "field1"
      },
      {
        "name": "field2",
        "type": "VARCHAR",
        "mapping": "field2"
      },
      {
        "name": "field3",
        "type": "BOOLEAN",
        "mapping": "field3"
      }
    ]
  }
}

.avsc 스키마 예시:

{
  "type" : "record",
  "name" : "example_avro_message",
  "namespace" : "io.trino.plugin.kafka",
  "fields" :
  [
    {
      "name":"field1",
      "type":["null", "long"],
      "default": null
    },
    {
      "name": "field2",
      "type":["null", "string"],
      "default": null
    },
    {
      "name":"field3",
      "type":["null", "boolean"],
      "default": null
    }
  ],
  "doc:" : "A basic avro schema"
}

protobuf 포맷

protobuf 포맷은 별도의 .proto 스키마 파일을 참조해요:

{
  "tableName": "example_table_name",
  "schemaName": "example_schema_name",
  "topicName": "example_topic_name",
  "key": { "..." },
  "message":
  {
    "dataFormat": "protobuf",
    "dataSchema": "/message_schema.proto",
    "fields":
    [
      {
        "name": "field1",
        "type": "BIGINT",
        "mapping": "field1"
      },
      {
        "name": "field2",
        "type": "VARCHAR",
        "mapping": "field2"
      },
      {
        "name": "field3",
        "type": "BOOLEAN",
        "mapping": "field3"
      }
    ]
  }
}

.proto 스키마 예시:

syntax = "proto3";

message schema {
  uint64 field1 = 1 ;
  string field2 = 2;
  bool field3 = 3;
}
INSERT INTO example_protobuf_table (field1, field2, field3)
  VALUES (123456789, 'example text', FALSE);

Protobuf oneof 지원:

syntax = "proto3";

message schema {
    oneof test_oneof_column {
        string string_column = 1;
        uint32 integer_column = 2;
        uint64 long_column = 3;
        double double_column = 4;
        float float_column = 5;
        bool boolean_column = 6;
    }
}

google.protobuf.Any 지원:

syntax = "proto3";

import "google/protobuf/any.proto";

message schema {
    google.protobuf.Any any_message = 1;
}

Any 메시지는 @type로 스키마를 가리켜야 해요:

{
  "@type":"file:///path/to/schemas/MyMessage",
  "longColumn":"493857959588286460",
  "numberColumn":"ONE",
  "stringColumn":"Trino"
}

관련 스키마:

syntax = "proto3";

message MyMessage {
  string stringColumn = 1;
  uint32 integerColumn = 2;
  uint64 longColumn = 3;
}

더 알아보기 (Learn more)

Kafka 커넥터를 실습으로 배우고 싶다면 Kafka 커넥터 튜토리얼 문서를 이어서 읽어 보세요.