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-provider와kafka-protobuf-typesJAR 파일을 복사해 클러스터의 모든 노드에서 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 커넥터 튜토리얼 문서를 이어서 읽어 보세요.