メインコンテンツまでスキップ

Kafka コネクタ リファレンス

このページでは、Lakeflow Connectにおけるマネージド Apache Kafka コネクタのコネクタオプション、テーブル構成オプション、および JSON トランスフォーマー設定について説明します。

コネクターオプション​

次のオプションは、取り込みパイプラインの各宛先テーブルに対してKafkaソースを構成します。パイプライン定義の connector_options.kafka_options で、これらのオプションを指定してください。完全なパイプラインの例については、「例」を参照してください。

オプション

Type

デフォルト

説明

topics

文字列のリスト:

—

サブスクライブするトピック名のリストです。topic_patternと相互に排他的です。topicsまたはtopic_patternのいずれかが必要です。

topic_pattern

String

—

サブスクライブするトピック名に一致する Java 正規表現です。topicsと相互に排他的です。

starting_offset

String

latest

チェックポイントが存在しない場合に読み取りを開始する場所(初回実行のみ)。有効な値: latest、earliest。

key_transformer

Transformer

—

メッセージキーのデシリアライザ構成。設定されていない場合、キー列は BINARY として保持されます。Transformer オプションを参照してください。

value_transformer

Transformer

—

メッセージ値の逆シリアライザー設定。設定されていない場合、値の列はBINARYとして保持されます。Transformerオプションを参照してください。

オプション

Type

デフォルト

説明

topics

文字列のリスト:

—

サブスクライブするトピック名のリストです。topic_patternと相互に排他的です。topicsまたはtopic_patternのいずれかが必要です。

topic_pattern

String

—

サブスクライブするトピック名に一致する Java 正規表現です。topicsと相互に排他的です。

starting_offset

String

latest

チェックポイントが存在しない場合に読み取りを開始する場所(初回実行のみ)。有効な値: latest、earliest。

key_transformer

Transformer

—

メッセージキーのデシリアライザ構成。設定されていない場合、キー列は BINARY として保持されます。Transformer オプションを参照してください。

value_transformer

Transformer

—

メッセージ値の逆シリアライザー設定。設定されていない場合、値の列はBINARYとして保持されます。Transformerオプションを参照してください。

Transformer オプション​

トランスフォーマーは、バイナリのKafkaメッセージのキーと値が構造化された列にどのように逆シリアル化されるかを定義します。シリアル化形式と、key_transformerまたはvalue_transformerの下の対応する形式固有のオプションを指定します。キー、値、またはその両方に対してトランスフォーマーを個別に構成できます。トランスフォーマーが設定されていない場合、列はBINARYとして保持されます。

JSON の場合、明示的なスキーマを提供するか、進化を含むスキーマ推論を使用するか、json_optionsを完全に省略して値をVARIANT列として保存できます。

オプション

適用対象

Type

デフォルト

説明

format

すべて

String

—

データのシリアル化形式。有効な値: STRING、JSON、AVRO、PROTOBUF。STRING に追加のオプションは必要ありません。トランスフォーマーで json_optionsが指定されていない場合、値は default として解析されます。VARIANT詳細については、バリアント データのフォーマットを参照してください。AVRO および PROTOBUF については、Avro オプションおよびProtobuf オプションを参照してください。

json_options.schema

JSON

String

—

Spark DDL形式のインラインスキーマ(例: "id BIGINT, name STRING")。schema_file_pathと相互に排他的です。

json_options.schema_file_path

JSON

String

—

.ddlスキーマファイルへのパス。schemaと相互に排他的です。Unity Catalog ボリュームパス (/Volumes/...) をサポートしています。

json_options.schema_evolution_mode

JSON

String

—

自動スキーマ推論のスキーマ進化モード。スキーマ進化モードを参照してください。

json_options.schema_hints

JSON

String

—

スキーマ推論に影響を与えるコンマ区切りの "column_name type" ペア (例: "id BIGINT, ts TIMESTAMP")。schema_evolution_mode が設定されている必要があります。スキーマヒントを使用したスキーマ推論のオーバーライドを参照してください。

オプション

適用対象

Type

デフォルト

説明

format

すべて

String

—

データのシリアル化形式。有効な値: STRING、JSON、AVRO、PROTOBUF。STRING に追加のオプションは必要ありません。トランスフォーマーで json_optionsが指定されていない場合、値は default として解析されます。VARIANT詳細については、バリアント データのフォーマットを参照してください。AVRO および PROTOBUF については、Avro オプションおよびProtobuf オプションを参照してください。

json_options.schema

JSON

String

—

Spark DDL形式のインラインスキーマ(例: "id BIGINT, name STRING")。schema_file_pathと相互に排他的です。

json_options.schema_file_path

JSON

String

—

.ddlスキーマファイルへのパス。schemaと相互に排他的です。Unity Catalog ボリュームパス (/Volumes/...) をサポートしています。

json_options.schema_evolution_mode

JSON

String

—

自動スキーマ推論のスキーマ進化モード。スキーマ進化モードを参照してください。

json_options.schema_hints

JSON

String

—

スキーマ推論に影響を与えるコンマ区切りの "column_name type" ペア (例: "id BIGINT, ts TIMESTAMP")。schema_evolution_mode が設定されている必要があります。スキーマヒントを使用したスキーマ推論のオーバーライドを参照してください。

Avro オプション​

format: AVRO の場合に、avro_options の下にこれらのオプションを設定します。スキーマは、インライン、ファイル、またはスキーマレジストリから指定してください。

オプション

Type

デフォルト

説明

avro_options.schema

String

—

JSON 形式のインライン Avro スキーマ。schema_file_path および schema_registry とは排他的(同時に使用不可)です。

avro_options.schema_file_path

String

—

.avscスキーマファイルへのパス。Unity Catalog Volumes のパス (/Volumes/...) をサポートします。schema および schema_registry と相互に排他的です。

avro_options.schema_registry

オブジェクト

—

schema または schema_file_path の代わりに、ランタイム時にスキーマレジストリからスキーマを解決します。スキーマレジストリのオプションを参照してください。

avro_options.parse_mode

String

PERMISSIVE

逆シリアル化に失敗したレコードの処理方法。有効な値: PERMISSIVE (default)。逆シリアル化に失敗した各レコードの未加工バイトを _corrupt_record 列(タイプ BINARY)に書き込み、レコードの他の列を null に設定して、処理を継続します。また、FAILFAST は、逆シリアル化に失敗した最初のレコードでパイプラインを失敗させます。

オプション

Type

デフォルト

説明

avro_options.schema

String

—

JSON 形式のインライン Avro スキーマ。schema_file_path および schema_registry とは排他的(同時に使用不可)です。

avro_options.schema_file_path

String

—

.avscスキーマファイルへのパス。Unity Catalog Volumes のパス (/Volumes/...) をサポートします。schema および schema_registry と相互に排他的です。

avro_options.schema_registry

オブジェクト

—

schema または schema_file_path の代わりに、ランタイム時にスキーマレジストリからスキーマを解決します。スキーマレジストリのオプションを参照してください。

avro_options.parse_mode

String

PERMISSIVE

逆シリアル化に失敗したレコードの処理方法。有効な値: PERMISSIVE (default)。逆シリアル化に失敗した各レコードの未加工バイトを _corrupt_record 列(タイプ BINARY)に書き込み、レコードの他の列を null に設定して、処理を継続します。また、FAILFAST は、逆シリアル化に失敗した最初のレコードでパイプラインを失敗させます。

Protobuf オプション​

format: PROTOBUFのときに、これらのオプションをprotobuf_optionsの下に設定します。コンパイル済みの記述子セット(.desc)ファイルとメッセージ名を指定するか、スキーマレジストリからスキーマを解決します。

オプション

Type

デフォルト

説明

protobuf_options.desc_file_path

String

—

コンパイル済みの Protobuf 記述子セット(.desc)ファイルへのパス。Unity Catalog ボリュームのパス(/Volumes/...)をサポートしています。schema_registryが設定されていない限り必須です。

protobuf_options.message_name

String

—

完全修飾されたProtobufメッセージ型の名前(例:com.example.events.UserEvent)。desc_file_path で必須です。

protobuf_options.schema_registry

オブジェクト

—

desc_file_path ではなく、ランタイム時にスキーマレジストリからスキーマを解決します。スキーマレジストリのオプションを参照してください。

protobuf_options.recursive_fields_max_depth

Integer

—

再帰的な Protobuf フィールドの最大展開深度。Spark SQL ではネイティブでサポートされていません。有効な値:-1(再帰的フィールドを禁止)、0(破棄)、1~10。

protobuf_options.parse_mode

String

PERMISSIVE

逆シリアル化に失敗したレコードの処理方法。有効な値: PERMISSIVE (default)。逆シリアル化に失敗した各レコードの未加工バイトを _corrupt_record 列(タイプ BINARY)に書き込み、レコードの他の列を null に設定して、処理を継続します。また、FAILFAST は、逆シリアル化に失敗した最初のレコードでパイプラインを失敗させます。

オプション

Type

デフォルト

説明

protobuf_options.desc_file_path

String

—

コンパイル済みの Protobuf 記述子セット(.desc)ファイルへのパス。Unity Catalog ボリュームのパス(/Volumes/...)をサポートしています。schema_registryが設定されていない限り必須です。

protobuf_options.message_name

String

—

完全修飾されたProtobufメッセージ型の名前(例:com.example.events.UserEvent)。desc_file_path で必須です。

protobuf_options.schema_registry

オブジェクト

—

desc_file_path ではなく、ランタイム時にスキーマレジストリからスキーマを解決します。スキーマレジストリのオプションを参照してください。

protobuf_options.recursive_fields_max_depth

Integer

—

再帰的な Protobuf フィールドの最大展開深度。Spark SQL ではネイティブでサポートされていません。有効な値:-1(再帰的フィールドを禁止)、0(破棄)、1~10。

protobuf_options.parse_mode

String

PERMISSIVE

逆シリアル化に失敗したレコードの処理方法。有効な値: PERMISSIVE (default)。逆シリアル化に失敗した各レコードの未加工バイトを _corrupt_record 列(タイプ BINARY)に書き込み、レコードの他の列を null に設定して、処理を継続します。また、FAILFAST は、逆シリアル化に失敗した最初のレコードでパイプラインを失敗させます。

Schema registry options​

Confluent互換のスキーマレジストリから実行時にスキーマを解決するには、avro_optionsまたはprotobuf_optionsの下でschema_registryを設定します。defaultで、パイプラインは、レジストリのURLとAPIキーを保存するパイプラインのKafkaソース接続を使用して、レジストリに対して認証します。別のUnity Catalog接続で認証するには、connection_nameを設定します。「接続プロパティ」を参照してください。

オプション

Type

デフォルト

説明

schema_registry.confluent_options.subject

String

—

必須。Confluent 互換のスキーマレジストリで解決するサブジェクト。

schema_registry.connection_name

String

—

レジストリに対して認証を行うために使用されるUnity Catalog接続。defaultでは、パイプラインのKafkaソース接続が使用されます。レジストリが異なる資格情報を使用する場合は、これを設定します。

schema_registry.protobuf_message_name

String

—

Protobuf のみ。サブジェクトが複数の Protobuf メッセージを定義している場合にメッセージを選択します。シンプル (Location) または完全修飾 (com.example.protos.Location)。default to the first message in the schema.

オプション

Type

デフォルト

説明

schema_registry.confluent_options.subject

String

—

必須。Confluent 互換のスキーマレジストリで解決するサブジェクト。

schema_registry.connection_name

String

—

レジストリに対して認証を行うために使用されるUnity Catalog接続。defaultでは、パイプラインのKafkaソース接続が使用されます。レジストリが異なる資格情報を使用する場合は、これを設定します。

schema_registry.protobuf_message_name

String

—

Protobuf のみ。サブジェクトが複数の Protobuf メッセージを定義している場合にメッセージを選択します。シンプル (Location) または完全修飾 (com.example.protos.Location)。default to the first message in the schema.

テーブル構成オプション​

以下のオプションは、connector_options の兄弟要素であるテーブルオブジェクトの table_configuration の下に設定されます。パイプラインの完全な例については、「例」を参照してください。

オプション

Type

デフォルト

説明

source_metadata_column

String

—

各レコードの Kafka ソースメタデータを保持するために送信先テーブルに追加される構造体列の名前。ソースメタデータ列を参照してください。名前は key または value にすることはできません。

オプション

Type

デフォルト

説明

source_metadata_column

String

—

各レコードの Kafka ソースメタデータを保持するために送信先テーブルに追加される構造体列の名前。ソースメタデータ列を参照してください。名前は key または value にすることはできません。

ファンアウト オプション​

備考

プレビュー

この機能はプライベート プレビュー段階です。試用については、Databricksの連絡先にお問い合わせください。

ファンアウトオプションは、単一のKafkaソースからの各レコードを多数の宛先テーブルのいずれかにルーティングします。これらのオプションは、パイプライン定義内のスキーマオブジェクト (テーブルオブジェクトではない) の fanout_options で指定します。完全なパイプラインの例については、複数テーブルへのレコードのルーティング (ファンアウト) を、制約については、ファンアウトの制限 を参照してください。

オプション

Type

デフォルト

説明

fanout_by

String

—

必須。生のソースレコード(Kafka の key および value 列を含む)に対して評価され、その結果が宛先テーブル名を決定する SQL 式。 その値はテーブル名の最後のセグメントになります:{destination_catalog}.{destination_schema}.{value}。value 列はバイナリであるため、フィールドを抽出する前に、たとえば cast(value as string):event_type::string のように文字列にキャストします。式はNULL以外のSTRINGに解決され、その文字列は引用符やサニタイズなしでテーブル名セグメントとしてそのまま使用されるため、有効な引用符なしのテーブル識別子である必要があります。スペース、ドット、または引用符なしの識別子では無効なその他の文字を含む値は、書き込みに失敗します。完全に数字である値(たとえば、123)も同様です。値に文字またはアンダースコアも含まれている場合、先頭の数字は許可されます(例: 2024_events)。宛先テーブルは、まだ存在しない場合に自動的に作成されます。

transforms

Transformerのリスト

—

ルーティング後、宛先テーブルへの書き込み前に、ルーティングされた各レコードに適用される変換。ルーティング後に実行されるため、fanout_byが見る値には影響しません。最大1つの変換が許可されており、format: JSONを使用する必要があります。ファンアウト変換オプションを参照してください。

オプション

Type

デフォルト

説明

fanout_by

String

—

必須。生のソースレコード(Kafka の key および value 列を含む)に対して評価され、その結果が宛先テーブル名を決定する SQL 式。 その値はテーブル名の最後のセグメントになります:{destination_catalog}.{destination_schema}.{value}。value 列はバイナリであるため、フィールドを抽出する前に、たとえば cast(value as string):event_type::string のように文字列にキャストします。式はNULL以外のSTRINGに解決され、その文字列は引用符やサニタイズなしでテーブル名セグメントとしてそのまま使用されるため、有効な引用符なしのテーブル識別子である必要があります。スペース、ドット、または引用符なしの識別子では無効なその他の文字を含む値は、書き込みに失敗します。完全に数字である値(たとえば、123)も同様です。値に文字またはアンダースコアも含まれている場合、先頭の数字は許可されます(例: 2024_events)。宛先テーブルは、まだ存在しない場合に自動的に作成されます。

transforms

Transformerのリスト

—

ルーティング後、宛先テーブルへの書き込み前に、ルーティングされた各レコードに適用される変換。ルーティング後に実行されるため、fanout_byが見る値には影響しません。最大1つの変換が許可されており、format: JSONを使用する必要があります。ファンアウト変換オプションを参照してください。

ファンアウト変換オプション​

transforms 内の各エントリでは、以下のオプションを使用します。ファンアウト変換では JSON 形式のみがサポートされています。

オプション

Type

デフォルト

説明

format

String

JSON

オプション。変換のシリアル化形式。ファンアウト変換でサポートされているのはJSONのみであり、省略された場合のdefaultです。STRING、AVRO、およびPROTOBUFはサポートされていません。

input_column

String

—

変換が読み取りおよび書き込みを行う列 (例: value)。このファンアウト JSON 変換は、この列をその場で解析します。

オプション

Type

デフォルト

説明

format

String

JSON

オプション。変換のシリアル化形式。ファンアウト変換でサポートされているのはJSONのみであり、省略された場合のdefaultです。STRING、AVRO、およびPROTOBUFはサポートされていません。

input_column

String

—

変換が読み取りおよび書き込みを行う列 (例: value)。このファンアウト JSON 変換は、この列をその場で解析します。

注記

output_column変換オプションは、ファンアウト変換には適用されません。ファンアウトJSON変換は、結果を常にその場でinput_columnに書き込みます。任意のoutput_columnの値は無視されます。

接続プロパティ​

Catalog Explorer で Unity Catalog Kafka 接続を作成する際、認証方法に応じて次のプロパティを指定する必要があります。接続作成ステップについては、「Kafka接続の作成」を参照してください。

ユーザー名とパスワード​

属性

説明

接続名

Unity Catalog内の接続の一意の名前。

接続タイプ

Select Kafka .

認証タイプ

ユーザー名とパスワード を選択します。

ユーザー名

Kafka クラスターでの認証に使用されるユーザー名。

パスワード

Kafka クラスターでの認証に使用されるパスワード。

ブートストラップサーバー

Kafkaクラスターのブートストラップサーバーアドレス (例: broker1:9092,broker2:9092)。

**スキーマレジストリURL**(オプション)

スキーマレジストリのURLです。

スキーマレジストリ API キー (オプション)

スキーマレジストリの API キーです。

スキーマレジストリ API シークレット(オプション)

スキーマレジストリの API シークレット。

属性

説明

接続名

Unity Catalog内の接続の一意の名前。

接続タイプ

Select Kafka .

認証タイプ

ユーザー名とパスワード を選択します。

ユーザー名

Kafka クラスターでの認証に使用されるユーザー名。

パスワード

Kafka クラスターでの認証に使用されるパスワード。

ブートストラップサーバー

Kafkaクラスターのブートストラップサーバーアドレス (例: broker1:9092,broker2:9092)。

**スキーマレジストリURL**(オプション)

スキーマレジストリのURLです。

スキーマレジストリ API キー (オプション)

スキーマレジストリの API キーです。

スキーマレジストリ API シークレット(オプション)

スキーマレジストリの API シークレット。

サービス資格情報​

属性

説明

接続名

Unity Catalog内の接続の一意の名前。

接続タイプ

Select Kafka .

認証タイプ

「**サービス資格情報**」を選択します。

サービス認証情報

既存の Unity Catalog サービス認証情報を選択するか、 新しいサービス認証情報を作成 をクリックします。

ブートストラップサーバー

Kafkaクラスターのブートストラップサーバーアドレス (例: broker1:9092,broker2:9092)。

**スキーマレジストリURL**(オプション)

スキーマレジストリのURLです。

スキーマレジストリ API キー (オプション)

スキーマレジストリの API キーです。

スキーマレジストリ API シークレット(オプション)

スキーマレジストリの API シークレット。

属性

説明

接続名

Unity Catalog内の接続の一意の名前。

接続タイプ

Select Kafka .

認証タイプ

「**サービス資格情報**」を選択します。

サービス認証情報

既存の Unity Catalog サービス認証情報を選択するか、 新しいサービス認証情報を作成 をクリックします。

ブートストラップサーバー

Kafkaクラスターのブートストラップサーバーアドレス (例: broker1:9092,broker2:9092)。

**スキーマレジストリURL**(オプション)

スキーマレジストリのURLです。

スキーマレジストリ API キー (オプション)

スキーマレジストリの API キーです。

スキーマレジストリ API シークレット(オプション)

スキーマレジストリの API シークレット。

宛先テーブルスキーマ​

Kafkaコネクタはストリーミングテーブル(追記専用)に書き込みます。宛先テーブルに書き込まれる列は、トランスフォーマーが構成されているかどうかによって異なります。

トランスフォーマーなし (生バイナリ)​

key_transformerまたはvalue_transformerが設定されていない場合、宛先テーブルには以下の列が含まれます。

列

Type

説明

key

BINARY

Kafkaメッセージキーの生バイナリコンテンツ。

value

BINARY

Kafkaメッセージ値の生のバイナリコンテンツ。

列

Type

説明

key

BINARY

Kafkaメッセージキーの生バイナリコンテンツ。

value

BINARY

Kafkaメッセージ値の生のバイナリコンテンツ。

文字列トランスフォーマーを使用​

トランスフォーマーでformat: STRINGが設定されている場合、対応する列はBINARYの代わりにSTRINGとして書き込まれます。

JSONトランスフォーマーを使用すると​

トランスフォーマーで format: JSON が設定されている場合:

  • json_optionsが指定されていない場合、列はVARIANTとして書き込まれます。
  • json_options.schemaまたはjson_options.schema_file_pathが指定されている場合、JSONはスキーマと一致する型付き列に解析されます。
  • json_options.schema_evolution_modeが設定されている場合、スキーマ推論が使用され、スキーマは自動的に進化します。

Avro トランスフォーマーを使用する場合​

トランスフォーマーで format: AVRO が設定されている場合、列は Avro スキーマに一致する型付きの列に解析されます。PERMISSIVE モード(default)では、型 BINARY の _corrupt_record 列も追加されます。この列には、逆シリアル化に失敗したすべてのレコードの未処理バイトが格納され、正常に解析されたレコードの場合は null になります。

Protobuf トランスフォーマーを使用する場合​

トランスフォーマーで format: PROTOBUF が設定されている場合、列は Protobuf メッセージ定義に一致する型付き列に解析されます。PERMISSIVE モード (default) では、型 BINARY の _corrupt_record 列も追加されます。これには、デシリアライズに失敗したレコードの未加工バイトが保持され、正常に解析されたレコードには null が設定されます。

ソースメタデータ列​

table_configuration の下で source_metadata_column を設定し、その名前の構造体列を宛先テーブルに追加します。この構造体には、各レコードの以下の Kafka ソースメタデータフィールドが含まれています:topic、partition、offset、timestamp、timestampType、および headers。列名には key または value を使用できません。テーブル構成オプションを参照してください。