AI エージェント向け: ドキュメントインデックスは https://www.mongodb.com/ja-jp/docs/llms.txt で利用できます。すべてのページの markdown バージョンは、いずれかの URL パスに .md を追加することで利用できます。
See how MongoDB 9.0 delivers up to 2x higher throughput.
MongoDB Branding Shape
Register now >
Docs Menu

$iceberg 集計ステージ

ステージでは、$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ステージは、次のフィールドを持つドキュメントを取得します。

フィールド
タイプ
必要性
説明

connectionName

string

必須

読み取りと書込みに使用するAWS S3 接続の名前。これは、接続レジストリ内の接続の名前と一致する必要があります。

bucket

string

必須

ターゲットApache Opsageデータベースを含む S3 バケットの名前。

databaseName

string

必須

ターゲット テーブルを含む Apache Iceberg データベースの名前。

tableName

string |式

必須

ターゲット Apache Iceberg テーブルの名前。string または string に評価される式である必要があります。ドキュメントごとのダイナミックルーティングに式を使用します。

path

string

必須

Apache Iceberg データベースへのパスのプレフィックス キー。

region

string

条件付き

バケットのAWSリージョン。 AWSで実行中いないストリーム プロセッサに必要です。

mode

string

任意

入力ドキュメントごとに実行する操作を決定する戦略。

  • "cdc" Atlas Stream Processing は、stream.source.operationType メタデータフィールドを読み取ることで操作タイプを決定します。

  • "insert" Atlas Stream Processingにより、stream.source.operationTypeメタデータフィールドの操作タイプ宣言を無視して、各ドキュメントがターゲットテーブルに新しい行として追加されます。

デフォルトは cdc です。

idFieldName

string

任意

cdc モードでローキーとして使用されるフィールド名と列名。

デフォルトは "_id" です。

partitionedBy

ドキュメント

任意

パーティショニング仕様。このフィールドを設定しない場合、$iceberg は idFieldName 列のデフォルトのパーティション変換を設定します。

1 つ以上のキーと値のペアを含むドキュメントである必要があります。各キーは、パーティション変換を実行する列の名前であり、各値は使用するパーティション変換である必要があります。最初のパーティション変換は idFieldName に対するものである必要があります。

struct 列内にネストされたフィールドでパーティション分割するには、フィールドへのドット区切りパス(例: outer.inner.key_field)を指定します。パスにはどのレベルでも list 列を含めることはできません。

所定のフィールドのパーティション変換値は、次のいずれかである必要があります。

  • "identity"

  • "year"

  • "month"

  • "day"

  • "hour"

  • { truncate: int }

  • { bucket: int }

Apache Iceberg パーティション変換の詳細については、Apache Iceberg ドキュメント を参照してください。

catalog

ドキュメント

任意

使用する Iceberg カタログ を定義するドキュメント。"hadoop" または "glue" の値を持つ type フィールドを含むドキュメントである必要があります。

schemaInference

ドキュメント

任意

$iceberg が テーブルスキーマを推論する方法を構成するドキュメント。次のいずれかの値を持つ modeフィールドを含むドキュメントである必要があります。

  • "json" により、Atlas Stream Processing は object フィールドと array フィールドをJSON文字列列として書込みます。

  • "nested" により、Atlas Stream Processing は object フィールドを struct 列として書込み、array フィールドを list 列として書込みます。

デフォルトは"json" です。詳細については、「 型の変換 」を参照してください。

$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 プリミティブ
詳細

string

string

int

int

long

long

double

double

bool

boolean

ObjectId

string

16 進数エンコード

UUID

string

文字列化されたUUID

BinData

binary

UUID には適用されません。

date

timestamptz

UTC 時間、マイクロ秒単位で測定

timestamp

timestamptz

UTC 時間、マイクロ秒まで測定

object

string or struct

デフォルトでは 基本JSON string としてシリアル化されます。nested モードでは、 はstruct 列になります。詳細については、「 オブジェクト フィールドと配列フィールド 」を参照してください。

array

string or list

デフォルトでは 基本JSON string としてシリアル化されます。nested モードでは、 はlist 列になります。詳細については、「 オブジェクト フィールドと配列フィールド 」を参照してください。

その他の BSON types はサポートされていません。Atlas Stream Processing は、サポートされていない BSON types のドキュメントをDLQ に送信します。

フィールドは、 schemaInference.modeAtlas Stream Processing が フィールドとobjectarray フィールドをApacheテーブルに書き込む方法を決定します。

  • jsonモードでは、Atlas Stream Processingobject arraystringjsonは フィールドと フィールドを基本JSON文字列として直列化し、 列に書込みます。 はデフォルトのモードです。

  • nestedモードでは、Atlas Stream Processing は フィールドと フィールドから Ops Managerobject arrayのネストされた型を推論します。

    • 各 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 ステージの様々なアプリケーションを示しています。

次の例は、Atlas データベースの初期コンテンツと変更ストリームを追加のみの形式で Apache Iceberg テーブルに書き込み、そのデータベースの操作履歴の耐久性のあるアーカイブを作成する方法を示しています。この集計には 2 つのステージがあります。

  1. $source ステージでは、Atlas データベースとの接続を確立し、特定の db データベース内の orders コレクションをターゲットにします。これにより、プロセッサーの有効化時にデータベース内のドキュメントをキャプチャする最初の同期が可能になり、各チェンジストリーム イベントでフル ドキュメントがキャプチャされることが保証されます。

  2. $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 つのステージがあります。

  1. $source ステージでは、Atlas データベースとの接続を確立し、特定の db データベース内の orders コレクションをターゲットにします。これにより、プロセッサーの有効化時にデータベース内のドキュメントをキャプチャする最初の同期が可能になり、各チェンジストリーム イベントでフル ドキュメントがキャプチャされることが保証されます。

  2. $match ステージはoperationTypeでフィルターし、有効な操作タイプ宣言を持つドキュメントのみを処理します。

  3. $replaceRoot ステージは、操作の種類に応じてドキュメントルートを変更します。

    • 削除操作の場合、ドキュメントのルートをドキュメントのキーに変更します。これにより、ドキュメントが削除されたことを示すレコードが作成されますが、そのコンテンツはそれ以降のプロセシングから除外されます。

    • その他のすべての操作については、ドキュメント ルートを fullDocument に変更し、ドキュメントのコンテンツをさらにプロセシングするために渡し、changestream メタデータを除外します。

  4. $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"
}
}

次の例では、Atlas Stream Processingは動的式を使用して、ドキュメントを様々な出力先に動的にルーティングします。

  1. $source ステージでは、Atlas データベースとの接続を確立し、特に db データベース内のコレクション a、b、cを対象とします。これにより、プロセッサーの有効化時にデータベース内のドキュメントをキャプチャする最初の同期が可能になり、各チェンジストリーム イベントでフル ドキュメントがキャプチャされることが保証されます。

  2. $match ステージは、operationType が "insert"、"update"、"delete"、"replace" のいずれかであるドキュメントをフィルターします。

  3. $replaceRoot ステージは、操作の種類に応じてドキュメントルートを変更します。

    • 削除操作の場合、ドキュメントのルートをドキュメントのキーに変更します。これにより、ドキュメントが削除されたことを示すレコードが作成されますが、そのコンテンツはそれ以降のプロセシングから除外されます。

    • その他のすべての操作については、ドキュメント ルートを fullDocument に変更し、ドキュメントのコンテンツをさらにプロセシングするために渡し、changestream メタデータを除外します。

  4. $icebergステージは という名前のAWS S3 バケットへの接続を確立し、myData iceberg-warehouse/パスのApacheテーブルに書き込みます。ドキュメントメタデータから検索されたソースcollection の名前に基づいて、テーブルの名前が決定されます。また、ドキュメントメタデータに従って実行する操作も決定します。

{
"$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"
}
}