定義
ステージでは、$iceberg 接続レジストリ でAWS S3 バケットへの接続を指定し、 Apache inserterger テーブルにデータを書込むことができます。
配置
$iceberg パイプラインの最後のステージである必要があります。1 つのパイプラインにつき、1 つの $iceberg ステージのみを使用できます。
構文
$icebergパイプライン ステージには次のプロトタイプ形式があります。
{ "$iceberg": { "connectionName": "<registered-connection>", "bucket": "<target-bucket>", "databaseName": "<database>", "tableName": "<string>" | <expression>, "path": "<key-prefix>", "region": "<target-region>", "mode": "cdc" | "insert", "idFieldName": "<field-name>", "partitionedBy": { "<column-name>": "<partition-transform>", . . . }, "catalog": { "type": "hadoop" | "glue" }, "schemaInference": { "mode": "json" | "nested" } } }
$icebergステージは、次のフィールドを持つドキュメントを取得します。
フィールド | タイプ | 必要性 | 説明 |
|---|---|---|---|
| string | 必須 | |
| string | 必須 | ターゲットApache Opsageデータベースを含む S3 バケットの名前。 |
| string | 必須 | ターゲット テーブルを含む Apache Iceberg データベースの名前。 |
| string |式 | 必須 | ターゲット Apache Iceberg テーブルの名前。string または string に評価される式である必要があります。ドキュメントごとのダイナミックルーティングに式を使用します。 |
| string | 必須 | Apache Iceberg データベースへのパスのプレフィックス キー。 |
| string | 条件付き | バケットのAWSリージョン。 AWSで実行中いないストリーム プロセッサに必要です。 |
| string | 任意 | 入力ドキュメントごとに実行する操作を決定する戦略。
デフォルトは |
| string | 任意 |
デフォルトは |
| ドキュメント | 任意 | パーティショニング仕様。このフィールドを設定しない場合、 1 つ以上のキーと値のペアを含むドキュメントである必要があります。各キーは、パーティション変換を実行する列の名前であり、各値は使用するパーティション変換である必要があります。最初のパーティション変換は
所定のフィールドのパーティション変換値は、次のいずれかである必要があります。
Apache Iceberg パーティション変換の詳細については、Apache Iceberg ドキュメント を参照してください。 |
| ドキュメント | 任意 | 使用する Iceberg カタログ を定義するドキュメント。 |
| ドキュメント | 任意 |
デフォルトは |
動作
$icebergステージを使用する場合、それはストリームプロセッサーの最後のステージである必要があります。
Atlas Stream Processing は、SP10、SP30、SP50 ストリーム プロセッサーのみ $iceberg ステージをサポートします。プロセッサー階層によって、ダイナミック ルーティングでサポートされるテーブルの最大数が決定されます。
階層 | 最大テーブル |
|---|---|
SP 10 | 5 |
SP 30 | 10 |
SP 50 | 50 |
重要
動的ルーティングを持つストリーム プロセッサが、その階層でサポートされているテーブルの最大数を超えると、プロセッサはFAILED 状態になります。障害の原因と回復の詳細については、「 エラー処理と再試行ポリシー 」を参照してください。
ステージは、ストリーム$iceberg プロセッサの出力データのスキーマから、結果として得られるApache insertock テーブルのスキーマを推測します。 Atlas Stream Processing がストリーム内の新しいフィールド(struct 列内のフィールドを含む)を観察すると、テーブルスキーマもそれに応じて変化します。
まだ存在しないテーブルを指定した場合、Apache Iceberg は、それを対象とする最初のメッセージを受信したときにテーブルを作成します。
Atlas Stream Processing は、Apache Iceberg テーブルへの出力に対して、少なくとも 1 回のプロセシングを保証します。
tableName フィールドの値として動的式を使用できます。動的式を使用してドキュメント固有の値を取得することで、これらの値に応じて入力ドキュメントを異なるテーブルにルーティングできます。式は string として評価される必要があります。例については、「動的ルーティング」を参照してください。詳細については、「式演算子」を参照してください。
動的な式でトピックを指定したものの、Atlas Stream Processing が特定のメッセージの式を評価できない場合、構成されている場合、Atlas Stream Processing はそのメッセージをデッドレターキュー(DLQ)に送信し、以降のメッセージを処理します。If there's no デッドレターキュー(DLQ) configured, then Atlas Stream Processing skips the message completely and processes subsequent messages.
型の変換
Atlas Stream Processing は $iceberg ステージのテーブルに書き込む際に、BSON から Iceberg primitive types への型変換を実行します。
BSON | Apache Iceberg プリミティブ | 詳細 |
|---|---|---|
|
| |
|
| |
|
| |
|
| |
|
| |
|
| 16 進数エンコード |
|
| 文字列化されたUUID |
|
| UUID には適用されません。 |
|
| UTC 時間、マイクロ秒単位で測定 |
|
| UTC 時間、マイクロ秒まで測定 |
|
| デフォルトでは 基本JSON string としてシリアル化されます。 |
|
| デフォルトでは 基本JSON string としてシリアル化されます。 |
その他の BSON types はサポートされていません。Atlas Stream Processing は、サポートされていない BSON types のドキュメントをDLQ に送信します。
オブジェクト フィールドと配列フィールド
フィールドは、 schemaInference.modeAtlas Stream Processing が フィールドとobjectarray フィールドをApacheテーブルに書き込む方法を決定します。
nestedモードでは、Atlas Stream Processing は フィールドと フィールドから Ops Managerobjectarrayのネストされた型を推論します。各
objectフィールドはstruct列になります。 Atlas Stream Processing は、オブジェクト内の対応するフィールドから、 構造体内の各フィールドの型を推論します。これは、ネストのすべてのレベルに適用されます。各
arrayフィールドはlist列になります。 Atlas Stream Processing は、配列の null 以外の要素からリスト要素の型を推論します。プリミティブ要素はすべて同じ型である必要があります。要素がオブジェクトの場合、要素型はすべてのオブジェクトのフィールドを含むstructです。
モードは、ターゲット テーブルに一致する列がないフィールドにのみ適用されます。列がすでに存在する場合、Atlas Stream Processing はモードを無視し、 列タイプ に基づいてフィールドを書込みます。
列が
stringの場合、Atlas Stream Processing はフィールドを基本JSON string として直列化します。列が
structまたはlistの場合、Atlas Stream Processing はそのネストされたタイプとしてフィールドを書込みます。
入力に新しい object または array フィールドが含まれている場合、Atlas Stream Processing はモードを使用して新しい列のタイプを推論します。
例
以下の例は、$iceberg ステージの様々なアプリケーションを示しています。
ChangeStream のアーカイブ
次の例は、Atlas データベースの初期コンテンツと変更ストリームを追加のみの形式で Apache Iceberg テーブルに書き込み、そのデータベースの操作履歴の耐久性のあるアーカイブを作成する方法を示しています。この集計には 2 つのステージがあります。
$sourceステージでは、Atlas データベースとの接続を確立し、特定のdbデータベース内のordersコレクションをターゲットにします。これにより、プロセッサーの有効化時にデータベース内のドキュメントをキャプチャする最初の同期が可能になり、各チェンジストリーム イベントでフル ドキュメントがキャプチャされることが保証されます。$icebergステージはAWS S3 バケットへの接続を確立し、myTableパス内のiceberg-warehouse/という名前のテーブルに書き込みます。insert操作のみを指定することで、追加専用のログスタイルの書込みフローが保証されます。
{ "$source": { "connectionName": "atlas1", "db": "db", "coll": "orders", "initialSync": { "enable": true }, "config": { "fullDocument": "required" } }, "$iceberg": { "connectionName": "myS3Connection", "bucket": "myData", "path": "iceberg-warehouse/", "tableName": "myTable", "mode": "insert" } }
コレクションのミラーリング
次の例は、Atlas コレクションを完全に Apache Iceberg テーブルにミラーリングする方法を示しています。
集計を定義する前に、次の変数を設定します。
const isDeleteExpr = {$eq: [{$meta: "stream.source.operationType"}, "delete"]};
次の集計では、Atlas ソース コレクションへの変更と同期して、Apache Iceberg テーブルエントリを追加、更新、削除します。これには 4 つのステージがあります。
$sourceステージでは、Atlas データベースとの接続を確立し、特定のdbデータベース内のordersコレクションをターゲットにします。これにより、プロセッサーの有効化時にデータベース内のドキュメントをキャプチャする最初の同期が可能になり、各チェンジストリーム イベントでフル ドキュメントがキャプチャされることが保証されます。$matchステージはoperationTypeでフィルターし、有効な操作タイプ宣言を持つドキュメントのみを処理します。$replaceRootステージは、操作の種類に応じてドキュメントルートを変更します。削除操作の場合、ドキュメントのルートをドキュメントのキーに変更します。これにより、ドキュメントが削除されたことを示すレコードが作成されますが、そのコンテンツはそれ以降のプロセシングから除外されます。
その他のすべての操作については、ドキュメント ルートを
fullDocumentに変更し、ドキュメントのコンテンツをさらにプロセシングするために渡し、changestream メタデータを除外します。
$icebergステージはAWS S3 バケットへの接続を確立し、myTableiceberg-warehouse/パスにある という名前のApache Tiger テーブルに書き込みます。cdcモードでは、このステージは各ドキュメントの メタデータフィールドから読み取ることで、 Apache Ticketrstream.source.operationTypeテーブルに対して実行する操作を決定します。
{ "$source": { "connectionName": "atlas1", "db": "db", "coll": "orders", "initialSync": { "enable": true }, "config": { "fullDocument": "required" } }, "$match": { "operationType": { "$in": ["insert", "update", "delete", "replace"] } }, "$replaceRoot": { "newRoot": { "$cond": { "if": isDeleteExpr, "then": "$documentKey", "else": "$fullDocument" } } } "$iceberg": { "connectionName": "myS3Connection", "bucket": "myData", "path": "iceberg-warehouse/", "tableName": "myTable", "mode": "cdc" } }
マルチコレクションソースからマルチテーブルApacheターゲットへ
次の例では、Atlas Stream Processingは動的式を使用して、ドキュメントを様々な出力先に動的にルーティングします。
$sourceステージでは、Atlas データベースとの接続を確立し、特にdbデータベース内のコレクションa、b、cを対象とします。これにより、プロセッサーの有効化時にデータベース内のドキュメントをキャプチャする最初の同期が可能になり、各チェンジストリーム イベントでフル ドキュメントがキャプチャされることが保証されます。$matchステージは、operationTypeが"insert"、"update"、"delete"、"replace"のいずれかであるドキュメントをフィルターします。$replaceRootステージは、操作の種類に応じてドキュメントルートを変更します。削除操作の場合、ドキュメントのルートをドキュメントのキーに変更します。これにより、ドキュメントが削除されたことを示すレコードが作成されますが、そのコンテンツはそれ以降のプロセシングから除外されます。
その他のすべての操作については、ドキュメント ルートを
fullDocumentに変更し、ドキュメントのコンテンツをさらにプロセシングするために渡し、changestream メタデータを除外します。
{ "$source": { "connectionName": "atlas1", "db": "db", "coll": ["a", "b", "c"], "initialSync": { "enable": true }, "config": { "fullDocument": "required" } } }, { "$match": { "operationType": { "$in": ["insert", "update", "delete", "replace"] } } }, { "$replaceRoot": { "newRoot": { "$cond": { "if": { "$eq": [{ "$meta": "stream.source.operationType" }, "delete"] }, "then": "$documentKey", "else": "$fullDocument" } } } }, { "$iceberg": { "connectionName": "myS3Connection", "databaseName": "iceberg-db", "bucket": "myData", "path": "iceberg-warehouse/", "tableName": { "$meta": "stream.source.ns.coll" }, "mode": "cdc" } }