メインコンテンツまでスキップ
バージョン: 4.x

Kafkaカタログ

概要

Kafka CatalogはTrino Kafka ConnectorをTrino Connector互換性フレームワーク経由で使用してKafka Topicデータにアクセスします。

注記
  • これは実験的機能で、バージョン3.0.1以降でサポートされています。
  • この機能はTrinoクラスタ環境に依存せず、Trino互換プラグインのみを使用します。

使用例

シナリオサポート状況
データ統合Kafka Topicデータを読み取りDoris内部テーブルに書き込み
データ書き戻しサポートされていません

バージョン互換性

  • Dorisバージョン: 3.0.1以降
  • Trino Connectorバージョン: 435
  • Kafkaバージョン: サポートされているバージョンについては、Trino Documentationを参照してください

クイックスタート

ステップ1: Connectorプラグインの準備

以下のいずれかの方法でKafka Connectorプラグインを取得できます:

方法1: プリコンパイル済みパッケージを使用(推奨)

プリコンパイル済みプラグインパッケージをこちらからダウンロードして展開してください。

方法2: 手動コンパイル

カスタムコンパイルが必要な場合は、以下の手順に従ってください(JDK 17が必要):

git clone https://github.com/apache/doris-thirdparty.git
cd doris-thirdparty
git checkout trino-435
cd plugin/trino-kafka
mvn clean package -Dmaven.test.skip=true

コンパイル後、trino/plugin/trino-kafka/target/の下にtrino-kafka-435/ディレクトリが作成されます。

Step 2: Deploy Plugin

  1. すべてのFEおよびBEデプロイメントパスのconnectors/ディレクトリにtrino-kafka-435/ディレクトリを配置します(存在しない場合は手動でディレクトリを作成してください):

    ├── bin
    ├── conf
    ├── plugins
    │ ├── connectors
    │ ├── trino-kafka-435
    ...

fe.conf 内の trino_connector_plugin_dir 設定を変更することで、プラグインパスをカスタマイズすることもできます。例:trino_connector_plugin_dir=/path/to/connectors/

  1. コネクタが適切に読み込まれるよう、すべてのFEおよびBEノードを再起動してください。

ステップ3:カタログの作成

基本設定

CREATE CATALOG kafka PROPERTIES (
'type' = 'trino-connector',
'trino.connector.name' = 'kafka',
'trino.kafka.nodes' = '<broker1>:<port1>,<broker2>:<port2>',
'trino.kafka.table-names' = 'test_db.topic_name',
'trino.kafka.hide-internal-columns' = 'false'
);

設定ファイルの使用

CREATE CATALOG kafka PROPERTIES (
'type' = 'trino-connector',
'trino.connector.name' = 'kafka',
'trino.kafka.nodes' = '<broker1>:<port1>,<broker2>:<port2>',
'trino.kafka.config.resources' = '/path/to/kafka-client.properties',
'trino.kafka.hide-internal-columns' = 'false'
);

デフォルトスキーマの設定

CREATE CATALOG kafka PROPERTIES (
'type' = 'trino-connector',
'trino.connector.name' = 'kafka',
'trino.kafka.nodes' = '<broker1>:<port1>,<broker2>:<port2>',
'trino.kafka.default-schema' = 'default_db',
'trino.kafka.hide-internal-columns' = 'false'
);

ステップ4: データクエリ

カタログを作成した後、3つの方法のいずれかを使用してKafka Topicデータをクエリできます:

-- Method 1: Switch to catalog and query
SWITCH kafka;
USE kafka_schema;
SELECT * FROM topic_name LIMIT 10;

-- Method 2: Use two-level path
USE kafka.kafka_schema;
SELECT * FROM topic_name LIMIT 10;

-- Method 3: Use fully qualified name
SELECT * FROM kafka.kafka_schema.topic_name LIMIT 10;

Schema Registry統合

Kafka CatalogはConfluent Schema Registryを通じた自動スキーマ取得をサポートしており、テーブル構造を手動で定義する必要がありません。

Schema Registryの設定

Basic認証

CREATE CATALOG kafka PROPERTIES (
'type' = 'trino-connector',
'trino.connector.name' = 'kafka',
'trino.kafka.nodes' = '<broker1>:<port1>',
'trino.kafka.table-description-supplier' = 'CONFLUENT',
'trino.kafka.confluent-schema-registry-url' = 'http://<schema-registry-host>:<schema-registry-port>',
'trino.kafka.confluent-schema-registry-auth-type' = 'BASIC_AUTH',
'trino.kafka.confluent-schema-registry.basic-auth.username' = 'admin',
'trino.kafka.confluent-schema-registry.basic-auth.password' = 'admin123',
'trino.kafka.hide-internal-columns' = 'false'
);

完全な設定例

CREATE CATALOG kafka PROPERTIES (
'type' = 'trino-connector',
'trino.connector.name' = 'kafka',
'trino.kafka.nodes' = '<broker1>:<port1>',
'trino.kafka.default-schema' = 'nrdp',
'trino.kafka.table-description-supplier' = 'CONFLUENT',
'trino.kafka.confluent-schema-registry-url' = 'http://<schema-registry-host>:<schema-registry-port>',
'trino.kafka.confluent-schema-registry-auth-type' = 'BASIC_AUTH',
'trino.kafka.confluent-schema-registry.basic-auth.username' = 'admin',
'trino.kafka.confluent-schema-registry.basic-auth.password' = 'admin123',
'trino.kafka.config.resources' = '/path/to/kafka-client.properties',
'trino.kafka.confluent-schema-registry-subject-mapping' = 'nrdp.topic1:NRDP.topic1',
'trino.kafka.hide-internal-columns' = 'false'
);

Schema Registry パラメータ

パラメータ名必須デフォルト説明
trino.kafka.table-description-supplierNo-Schema Registry サポートを有効にするには CONFLUENT に設定
trino.kafka.confluent-schema-registry-urlYes*-Schema Registry サービスアドレス
trino.kafka.confluent-schema-registry-auth-typeNoNONE認証タイプ:NONE、BASIC_AUTH、BEARER
trino.kafka.confluent-schema-registry.basic-auth.usernameNo-Basic Auth ユーザー名
trino.kafka.confluent-schema-registry.basic-auth.passwordNo-Basic Auth パスワード
trino.kafka.confluent-schema-registry-subject-mappingNo-Subject名のマッピング、形式:<db1>.<tbl1>:<topic_name1>,<db2>.<tbl2>:<topic_name2>
ヒント

Schema Registryを使用する場合、DorisはSchema RegistryからTopicスキーマ情報を自動的に取得するため、テーブル構造を手動で作成する必要がなくなります。

Subject マッピング

場合によっては、Schema Registryに登録されたSubject名がKafkaのTopic名と一致せず、データクエリが実行できないことがあります。このような場合は、confluent-schema-registry-subject-mappingを通じてマッピング関係を手動で指定する必要があります。

-- Map schema.topic to SCHEMA.topic Subject in Schema Registry
'trino.kafka.confluent-schema-registry-subject-mapping' = '<db1>.<tbl1>:<topic_name1>'

db1tbl1はDorisで表示される実際のDatabaseとTable名であり、topic_name1はKafkaでの実際のTopic名です(大文字小文字を区別)。

複数のマッピングはカンマで区切ることができます:

'trino.kafka.confluent-schema-registry-subject-mapping' = '<db1>.<tbl1>:<topic_name1>,<db2>.<tbl2>:<topic_name2>'

設定

カタログ設定パラメータ

Kafka Catalogを作成するための基本構文は以下の通りです:

CREATE CATALOG [IF NOT EXISTS] catalog_name PROPERTIES (
'type' = 'trino-connector', -- Required, fixed value
'trino.connector.name' = 'kafka', -- Required, fixed value
{TrinoProperties}, -- Trino Connector related properties
{CommonProperties} -- Common properties
);

TrinoProperties パラメータ

TrinoPropertiesは、trino.というプレフィックスが付いたTrino Kafka Connector固有のプロパティを設定するために使用されます。一般的なパラメータには以下が含まれます:

パラメータ名必須デフォルト説明
trino.kafka.nodesYes-Kafka Brokerノードアドレスリスト、形式:host1:port1,host2:port2
trino.kafka.table-namesNo-マップするTopicsのリスト、形式:schema.topic1,schema.topic2
trino.kafka.default-schemaNodefaultデフォルトスキーマ名
trino.kafka.hide-internal-columnsNotrueKafka内部カラム(_partition_id_partition_offsetなど)を非表示にするかどうか
trino.kafka.config.resourcesNo-Kafkaクライアント設定ファイルパス
trino.kafka.table-description-supplierNo-テーブル構造プロバイダ、Schema Registryを使用するにはCONFLUENTに設定
trino.kafka.confluent-schema-registry-urlNo-Schema Registryサービスアドレス

より多くのKafka Connector設定パラメータについては、Trino Official Documentationを参照してください。

CommonProperties パラメータ

CommonPropertiesは、メタデータ更新ポリシーや権限制御などの一般的なカタログプロパティを設定するために使用されます。詳細については、Catalog Overviewの「Common Properties」セクションを参照してください。

Kafkaクライアント設定

高度なKafkaクライアントパラメータ(セキュリティ認証、SSLなど)を設定する必要がある場合、設定ファイルを通じて指定できます。設定ファイル(例:kafka-client.properties)を作成してください:

# ============================================
# Kerberos/SASL Authentication Configuration
# ============================================
sasl.mechanism=GSSAPI
sasl.kerberos.service.name=kafka

# JAAS Configuration - Using keytab method
sasl.jaas.config=com.sun.security.auth.module.Krb5LoginModule required \
useKeyTab=true \
storeKey=true \
useTicketCache=false \
serviceName="kafka" \
keyTab="/opt/trino/security/keytabs/kafka.keytab" \
principal="kafka@EXAMPLE.COM";

# ============================================
# Avro Deserializer Configuration
# ============================================
key.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer
value.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer

カタログを作成する際に設定ファイルを指定します:

CREATE CATALOG kafka PROPERTIES (
'type' = 'trino-connector',
'trino.connector.name' = 'kafka',
'trino.kafka.nodes' = '<broker1>:<port1>',
'trino.kafka.config.resources' = '/path/to/kafka-client.properties'
);

データ型マッピング

Kafka Catalogを使用する際、データ型は以下のルールに従ってマッピングされます:

Kafka/Avro TypeTrino TypeDoris TypeNotes
booleanbooleanboolean
intintegerint
longbigintbigint
floatrealfloat
doubledoubledouble
bytesvarbinarystringHEX(col)関数を使用してクエリ
stringvarcharstring
arrayarrayarray
mapmapmap
recordrowstruct複雑なネストされた構造
enumvarcharstring
fixedvarbinarystring
null--
ヒント
  • bytes型の場合、16進数形式で表示するにはHEX()関数を使用します。
  • Kafka Catalogでサポートされるデータ型は、使用されるシリアライゼーション形式(JSON、Avro、Protobufなど)とSchema Registryの設定によって決まります。

Kafka内部カラム

Kafka ConnectorはKafkaメッセージのメタデータ情報にアクセスするための内部カラムを提供します:

Column NameTypeDescription
_partition_idbigintメッセージが配置されているパーティションID
_partition_offsetbigintパーティション内のメッセージオフセット
_message_timestamptimestampメッセージのタイムスタンプ
_keyvarcharメッセージキー
_key_corruptbooleanキーが破損しているかどうか
_key_lengthbigintキーのバイト長
_messagevarchar生のメッセージ内容
_message_corruptbooleanメッセージが破損しているかどうか
_message_lengthbigintメッセージのバイト長
_headersmapメッセージヘッダー情報

デフォルトでは、これらの内部カラムは非表示になっています。これらのカラムをクエリする必要がある場合は、カタログ作成時に設定してください:

'trino.kafka.hide-internal-columns' = 'false'

クエリの例:

SELECT 
_partition_id,
_partition_offset,
_message_timestamp,
*
FROM kafka.schema.topic_name
LIMIT 10;

制限事項

  1. 読み取り専用アクセス: Kafka Catalogはデータの読み取りのみをサポートし、書き込み操作(INSERT、UPDATE、DELETE)はサポートされていません。

  2. テーブル名の設定: Schema Registryを使用しない場合、trino.kafka.table-namesパラメータを通じてアクセスするTopicのリストを明示的に指定する必要があります。

  3. スキーマ定義:

    • Schema Registryを使用する場合、スキーマ情報はSchema Registryから自動的に取得されます。
    • Schema Registryを使用しない場合、テーブル定義を手動で作成するか、TrinoのTopic記述ファイルを使用する必要があります。
  4. データ形式: サポートされるデータ形式は、Topicが使用するシリアライゼーション方法(JSON、Avro、Protobufなど)に依存します。詳細については、Trino Official Documentationを参照してください。

  5. パフォーマンスに関する考慮事項:

    • Kafka CatalogはKafkaデータをリアルタイムで読み取るため、大量のデータをクエリするとパフォーマンスに影響を与える可能性があります。
    • LIMIT句や時間フィルタ条件を使用して、スキャンするデータ量を制限することを推奨します。

機能デバッグ

機能検証のためのKafka環境を迅速に構築するには、hereを参照してください。

参考資料