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

$source ステージ(Stream Processing)

$source

$sourceステージでは、データをストリーミングするための接続を接続レジストリで指定します。 次の接続タイプがサポートされています。

  • Apache Kafkaエージェント

  • MongoDB コレクションの変更ストリーム

  • MongoDB database 変更ストリーム

  • MongoDBクラスターの変更ストリーム

  • AWS Kinesisデータストリーム

  • ドキュメント配列

  • CRON スケジュール

Apache Kafkaプロバイダーからのストリーミングデータを使用する場合、$source ステージには次のプロトタイプ形式があります。

{
"$source": {
"connectionName": "<registered-connection>",
"topic" : ["<source-topic>", ...],
"timeField": {
$toDate | $dateFromString: <expression>
},
"partitionIdleTimeout": {
"size": <duration-number>,
"unit": "<duration-unit>"
},
"schemaRegistry": {
"connectionName": "<schema-registry-name>",
},
"config": {
"auto_offset_reset": "<start-event>",
"group_id": "<group-id>",
"keyFormat": "<deserialization-type>",
"keyFormatError": "<error-handling>"
},
}
}

$sourceステージは、次のフィールドを持つドキュメントを取得します。

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

connectionName

string

必須

データを取り込む接続レジストリ内の接続を識別するラベル。

topic

文字列または複数の文字列の配列

必須

メッセージをストリーミングする 1 つ以上の Apache Kafka トピックの名前。複数のトピックからのメッセージをストリーミングする場合は、配列で指定します。

timeField

ドキュメント

任意

受信メッセージの権限のあるタイムスタンプを定義するドキュメント。

timeFieldを使用する場合は、次のいずれかとして定義する必要があります。

  • ソース メッセージフィールドを引数として受け取る :式:$toDate

  • ソース メッセージフィールドを引数として受け取る :式:$dateFromString

timeFieldを宣言しない場合、Atlas Stream Processing は、ソースによって提供されたメッセージ タイムスタンプからタイムスタンプを作成します。

partitionIdleTimeout

ドキュメント

任意

証明機関の計算で無視される前に、パーティションがアイドル状態になることを許可する時間を指定するドキュメント。

このフィールドはデフォルトで無効になっています。アイドル状態で進まないパーティションを処理するには、このフィールドに値を設定します。

partitionIdleTimeout.size

integer

任意

パーティションのアイドル タイムアウトの期間を指定する数値。

partitionIdleTimeout.unit

string

任意

パーティション アイドル タイムアウトの期間の単位。

unitの値は次のいずれかになります。

  • "ms" (ミリ秒)

  • "second"

  • "minute"

  • "hour"

  • "day"

schemaRegistry

ドキュメント

任意

平均直列化されたソースからの読み取りをサポートするためにスキーマ レジストリの使用を有効にするドキュメント。

この機能を有効にするには、スキーマ Registry 接続を作成する必要があります。

schemaRegistry.connectionName

string

条件付き

Avro 非直列化に使用するスキーマ レジストリ接続の名前。

config

ドキュメント

任意

のさまざまなデフォルト値を上書きするフィールドを含むドキュメント。

config.auto_offset_reset

string

任意

Apache Kafkaソーストピック内のどのイベントで取り込みを開始するかを指定します。auto_offset_reset は次の値を取ります。

  • end、 latest 、またはlargest : 集計が初期化されたときに、トピックの最新のイベントから取り込みを開始します。

  • earliest、 beginning 、またはsmallest : トピック内の最も近いイベントから取り込みを開始します。

デフォルトは latest です。

config.group_id

string

任意

ストリーム プロセッサと関連付ける Kafka コンシューマー グループの ID。省略した場合、Atlas Stream Processing は、Stream Processing ワークスペースを次の形式の自動生成された ID に関連付けます。

asp-${streamProcessorId}-consumer

Atlas Stream Processing は、すべての永続ストリーム プロセッサに対してこのパラメーターの値を自動的に生成します。SP.process() で定義されたエフェメラル ストリーム プロセッサの場合、このパラメーターは手動で定義した場合にのみ設定されます。

config.enable_auto_commit

ブール値

条件付き

Kafkaプロバイダー パーティション オフセットのコミット ポリシーを決定するフラグ。Atlas Stream Processing は 2 つのコミット ポリシーをサポートしています。

  • このパラメータを true に設定すると、$source ステージが次の演算子にデータを渡すたびに、Atlas Stream Processing はオフセットをコミットします。

  • このパラメータを false に設定すると、Atlas Stream Processing がチェックポイントを取得するときに、ストリーム プロセッサはパーティション オフセットをコミットします。

SP.process() で定義されたエフェメラル ストリーム プロセッサの場合、group_id を設定しない限り、このパラメータはデフォルトで false に設定されます。それ以外の場合、デフォルトは true になります。

Kafka を$source として使用する場合のオフセットの詳細については、 Kafkaソースとコンシューマー グループ オフセット を参照してください。

config.keyFormat

string

任意

Apache Kafkaキー データを逆直列化するために使用されるデータ型。次のいずれかの値である必要があります。

  • "binData"

  • "string"

  • "json"

  • "int"

  • "long"

デフォルトは binData です。

config.keyFormatError

string

任意

Apache Kafkaキー データを逆直列化するときに発生したエラーの処理方法。次のいずれかの値である必要があります。

注意

Atlas Stream Processing では、ソース データ ストリーム内のドキュメントが有効なjsonまたはejsonである必要があります。 Atlas Stream Processing は、この要件を満たさないドキュメントをデッド レター キューに設定します(デッド レター キューを設定している場合)。

Atlas コレクションの変更ストリームを使用すると、アプリケーションは単一のコレクションにおけるリアルタイムデータの変更にアクセスできます。コレクションに対する変更ストリームを開く方法については、変更ストリームを参照してください。

変更ストリーム$sourceを使用する場合は、ソースクラスターで少なくとも24時間のoplog windowを構成してください。

変更ストリームを読み取るために、Atlas Stream Processing は oplog コレクションをスキャンします。その結果、ログに COLLSCAN 警告が表示されることがあります。これらの警告は通常の動作を示しており、エラーを示すものではありません。

config.fullDocumentまたはconfig.fullDocumentBeforeChangeをrequiredに設定する場合は、キャプチャー対象の書き込み操作を実行する前に、各コレクションでchangeStreamPreAndPostImagesを有効にしてください。書き込み発生時に機能が有効化されていなかった場合、またはpost-imageの有効期限が切れていた場合に、イベントでpost-imageを利用できないと、ストリームプロセッサーは失敗します。変更前ドキュメントおよび変更後ドキュメントを有効化する方法については、「ドキュメントのPre-ImageおよびPost-Imageを使用したChange Streams」を参照してください。

重要

変更ストリームを再開できるのは、再開トークンが識別する操作をoplogコレクションが保持している間のみです。アプリケーションが変更ストリームを再開できるようにするには、予想される最長時間よりも長い最小 Oplog Window を設定します。

Atlas コレクションの変更ストリームからのストリーミング データを操作する場合、 $sourceステージには次のプロトタイプ形式があります。

{
"$source": {
"connectionName": "<registered-connection>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"db" : "<source-db>",
"coll" : ["<source-coll>",...],
"initialSync": {
"enable": <boolean>,
"parallelism": <integer>
},
"readPreference": "<read-preference>",
"readPreferenceTags": [
{"<key>": "<value>"},
. . .
],
"config": {
"startAfter": <start-token> | "startAtOperationTime": <timestamp>,
"fullDocument": "<full-doc-condition>",
"fullDocumentOnly": <boolean>,
"fullDocumentBeforeChange": "<before-change-condition>",
"pipeline": [{
"<aggregation-stage>" : {
<stage-input>,
. . .
},
. . .
}],
"maxAwaitTimeMS": <time-ms>,
}
}
}

$sourceステージは、次のフィールドを持つドキュメントを取得します。

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

connectionName

string

条件付き

データを取り込む接続レジストリ内の接続を識別するラベル。

timeField

ドキュメント

任意

受信メッセージの権限のあるタイムスタンプを定義するドキュメント。

timeFieldを使用する場合は、次のいずれかとして定義する必要があります。

  • ソース メッセージ フィールドを引数として$toDate式

  • ソース メッセージ フィールドを引数として$dateFromString式。

timeFieldを宣言しない場合、Atlas Stream Processing は、ソースによって提供されたメッセージ タイムスタンプからタイムスタンプを作成します。

db

string

必須

connectionNameによって指定された Atlas インスタンスでホストされている MongoDB database の名前。 このデータベースの変更ストリームは、ストリーミング データソースとして機能します。

coll

文字列または複数の文字列の配列

必須

によって指定された Atlasインスタンスでホストされている 1 つ以上のMongoDBコレクションの名前。これらのコレクションの変更ストリームは、ストリーミングデータソースとして機能します。このフィールドを省略すると、ストリームconnectionName プロセッサはMongoDB Database Change Stream を使用します。

initialSync

ドキュメント

任意

initialSync の構成パラメータを含むドキュメント。

Atlas Stream Processing initialSync を使用すると、 changeEvent ドキュメントを挿入するのと同様に、Atlasコレクションに既存のドキュメントを取り込むことができます。 initialSync を有効にしている場合、ストリーム プロセッサを起動すると、まずコレクション内のすべての既存のドキュメントを取り込んで処理してから、新しい受信 changeEvent ドキュメントの取り込みと処理に進みます。 initialSync が完了すると、繰り返されません。

collが複数のコレクションを指定する場合、Atlas Stream Processing はリスト内のすべてのコレクションを同期します。マルチコレクションの同期順序、階層制限、および例については、「 マルチコレクションの最初の同期の構成 」を参照してください。

Atlas Stream Processing は、ストリーム プロセッサの起動時に同期するコレクションのセットを修正します。点の後に作成するコレクションは同期されません。いずれかのコレクションの同期に失敗すると、ストリーム プロセッサ全体が失敗します。

initialSync を有効にすると、パイプラインで $hoppingWindow、$sessionWindow、または $tumblingWindow ステージを使用できなくなります。

重要: initialSyncの考慮事項と制約については、「 制限 」を参照してください。

initialSync.enable

ブール値

条件付き

initialSync を有効にするかどうかを決定します。 initialSyncフィールドを宣言する場合は、このフィールドを に設定する必要があります。

initialSync.parallelism

integer

任意

initialSync操作を処理する並列処理のレベルを決定します。値を指定しない場合、デフォルトは 1 になります。

coll が複数のコレクションを指定している場合、この値は個々のコレクションではなく、コレクションリスト全体に適用されます。

initialSync は、Atlas Stream Processing が予測可能な順序で読み取れないタイプを含む、サポートされている _id タイプのコレクションにこの値を適用します。

設定できる最大値は、ストリームプロセッサーの階層によって異なります。詳しくは、Atlas Stream Processing 階層選択ガイドを参照してください。

各ストリーム プロセッサには、その階層によって決定される最大累積並列処理値があります。ストリーム プロセッサの累積並列処理は、次のように計算されます。

parallelism total - parallelized stages

parallelism total は $source、$lookup、$merge、$emit、$externalFunction ステージにおいて 1 を超えるすべての parallelism 値の合計であり、parallelized stages は parallelism の値が 1 より大きいこれらのステージの数です。

例、$source ステージが 4 の parallelism 値を設定し、$lookup ステージでは parallelism 値が設定されていない(デフォルトは 1)、かつ $merge ステージが parallelism に設定されている場合: 2 の } 値がある場合は 2 つの parallelized stages があり、ストリーム プロセッサの累積並列処理は (4 + 2) - 2 として計算されます。

ストリーム プロセッサがその階層の最大累積並列処理を超える場合、Atlas Stream Processing はエラーをスローし、目的のレベルの並列処理に必要な最小プロセッサ階層について提案します。エラーを解決するには、プロセッサをより高い階層に増やすアップするか、ステージの並列処理値を低くする必要があります。詳しくは、Stream Processingをご覧ください。

readPreference

string

任意

読み込み設定 (read preference)変更ストリームとinitialSync 操作の 。

デフォルトは primary です。

readPreferenceTags

配列

任意

読み込み設定 (read preference) タグ変更ストリームとinitialSync 操作の。

config

ドキュメント

任意

のさまざまなデフォルト値を上書きするフィールドを含むドキュメント。

config.startAfter

token

条件付き

ソースがレポートを開始する変更イベント。 これは再開トークンの形式をとります。

config.startAfterまたはconfig.startAtOperationTimeのいずれか 1 つだけを使用できます。

config.startAtOperationTime

タイムスタンプ | date

条件付き

ソースがレポートを開始するoptime 。

config.startAfterまたはconfig.startAtOperationTimeのいずれか 1 つだけを使用できます。

MongoDB 拡張 JSON の $date または $timestamp の値を受け付けます。

config.fullDocument

string

条件付き

変更ストリーム ソースが完全なドキュメントを返すか、更新が発生したときにのみ変更を返すかを制御する設定。 次のいずれかである必要があります。

  • default : update 操作の場合、完全なドキュメントを返しません。

  • updateLookup : 更新操作によって行われた変更に加えて、過半数がコミットした現在のバージョンの完全なドキュメントを返します。

  • required : 完全なドキュメントを返す必要があります。完全なドキュメントが利用できない場合、ストリームプロセッサは失敗します。これは削除操作には適用されません。この操作では、完全なドキュメントが利用できない場合に削除イベントが作成されます。

  • whenAvailable : 完全なドキュメントが利用可能になるたびに完全なドキュメントを返し、そうでない場合は変更を返します。

コレクション変更ストリームで またはrequired whenAvailableを使用するには、そのコレクションで変更ストリームの事前イメージと事後イメージを有効にする必要があります。

config.fullDocumentOnly

ブール値

条件付き

すべてのメタデータを含む変更イベント ドキュメント全体を返すか、 fullDocumentの内容のみを返すかを制御する設定。 trueに設定されている場合、ソースはfullDocumentの内容のみを返します。

がfullDocument requiredまたはwhenAvailable に設定されている場合は、そのコレクションで変更ストリームの事前イメージと事後イメージを有効にする必要があります。

config.fullDocumentBeforeChange

string

任意

変更ストリーム ソースに、出力に元の「変更前」状態の完全なドキュメントを含めるかどうかを指定します。 次のいずれかである必要があります。

  • off : fullDocumentBeforeChangeフィールドを省略します。

  • required : 状態が変更される前に、完全なドキュメントを返す必要があります。 状態が変更される前の完全なドキュメントが利用できない場合、ストリーム プロセッサは失敗します。

  • whenAvailable : 使用可能な場合は常に、変更前の状態で完全なドキュメントを返します。それ以外の場合は、 fullDocumentBeforeChangeフィールドを省略します。

fullDocumentBeforeChangeの値を指定しない場合、デフォルトはoffになります。

コレクションの変更ストリームでこのフィールドを使用するには、そのコレクションで変更ストリームの事前イメージと事後イメージを有効にする必要があります。

config.pipeline

ドキュメント

任意

集計パイプラインを指定し、変更ストリーム出力がパスされ、さらなるプロセシングが実行される前にフィルタリングします。このパイプラインは「変更ストリーム出力の修正」で説明されているパラメータに準拠する必要があります。

重要: 各変更イベントには wallTime フィールドと clusterTime フィールドが含まれます。$source 以降の Atlas Stream Processing ステージでは、プロセッサがこれらのフィールドを取り込んだため、受信することが想定されています。Change Stream データの適切なプロセシングを確保するために、$source.config.pipeline ではこれらのフィールドを変更しないでください。

config.maxAwaitTimeMS

integer

任意

空のバッチを返す前に、新しいデータ変更が変更ストリーム カーソルに報告されるまで待機する最大時間(ミリ秒)。

デフォルトは 1000 です。

Atlas データベースの変更ストリームを使用すると、アプリケーションは単一のデータベースでリアルタイムデータの変更にアクセスできます。データベースに対して変更ストリームを開く方法については、変更ストリームを参照してください。

変更ストリーム$sourceを使用する場合は、ソースクラスターで少なくとも24時間のoplog windowを構成してください。

変更ストリームを読み取るために、Atlas Stream Processing は oplog コレクションをスキャンします。その結果、ログに COLLSCAN 警告が表示されることがあります。これらの警告は通常の動作を示しており、エラーを示すものではありません。

config.fullDocumentまたはconfig.fullDocumentBeforeChangeをrequiredに設定する場合は、キャプチャー対象の書き込み操作を実行する前に、各コレクションでchangeStreamPreAndPostImagesを有効にしてください。書き込み発生時に機能が有効化されていなかった場合、またはpost-imageの有効期限が切れていた場合に、イベントでpost-imageを利用できないと、ストリームプロセッサーは失敗します。変更前ドキュメントおよび変更後ドキュメントを有効化する方法については、「ドキュメントのPre-ImageおよびPost-Imageを使用したChange Streams」を参照してください。

重要

変更ストリームを再開できるのは、再開トークンが識別する操作をoplogコレクションが保持している間のみです。アプリケーションが変更ストリームを再開できるようにするには、予想される最長時間よりも長い最小 Oplog Window を設定します。

Atlas データベース変更ストリームからのストリーミング データを操作する場合、 $sourceステージには次のプロトタイプ形式があります。

{
"$source": {
"connectionName": "<registered-connection>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"db" : "<source-db>",
"readPreference": "<read-preference>",
"readPreferenceTags": [
{"<key>": "<value>"},
. . .
],
"config": {
"startAfter": <start-token> | "startAtOperationTime": <timestamp>,
"fullDocument": "<full-doc-condition>",
"fullDocumentOnly": <boolean>,
"fullDocumentBeforeChange": "<before-change-condition>",
"pipeline": [{
"<aggregation-stage>" : {
<stage-input>,
. . .
},
. . .
}]
},
}
}

$sourceステージは、次のフィールドを持つドキュメントを取得します。

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

connectionName

string

条件付き

データを取り込む接続レジストリ内の接続を識別するラベル。

timeField

ドキュメント

任意

受信メッセージの権限のあるタイムスタンプを定義するドキュメント。

timeFieldを使用する場合は、次のいずれかとして定義する必要があります。

  • ソース メッセージ フィールドを引数として$toDate式

  • ソース メッセージ フィールドを引数として$dateFromString式。

timeFieldを宣言しない場合、Atlas Stream Processing は、ソースによって提供されたメッセージ タイムスタンプからタイムスタンプを作成します。

db

string

必須

connectionNameによって指定された Atlas インスタンスでホストされている MongoDB database の名前。 このデータベースの変更ストリームは、ストリーミング データソースとして機能します。

readPreference

string

任意

readPreferenceTags

配列

任意

config

ドキュメント

任意

のさまざまなデフォルト値を上書きするフィールドを含むドキュメント。

config.startAfter

token

条件付き

ソースがレポートを開始する変更イベント。 これは再開トークンの形式をとります。

config.startAfterまたはconfig.startAtOperationTimeのいずれか 1 つだけを使用できます。

config.startAtOperationTime

タイムスタンプ | date

条件付き

ソースがレポートを開始するoptime 。

config.startAfterまたはconfig.startAtOperationTimeのいずれか 1 つだけを使用できます。

MongoDB 拡張 JSON の $date または $timestamp の値を受け付けます。

config.fullDocument

string

条件付き

変更ストリーム ソースが完全なドキュメントを返すか、更新が発生したときにのみ変更を返すかを制御する設定。 次のいずれかである必要があります。

  • default : サーバーのデフォルトの動作を使用します。update 操作の場合、完全なドキュメントは返されません。

  • updateLookup : 更新操作によって行われた変更に加えて、過半数がコミットした現在のバージョンの完全なドキュメントを返します。

  • required : 完全なドキュメントを返す必要があります。完全なドキュメントが利用できない場合、ストリームプロセッサは失敗します。これは削除操作には適用されません。この操作では、完全なドキュメントが利用できない場合に削除イベントが作成されます。

  • whenAvailable : 完全なドキュメントが利用可能になるたびに完全なドキュメントを返し、そうでない場合は変更を返します。

fullDocumentの値を指定しない場合、デフォルトはdefaultになります。

データベース変更ストリームで またはrequired whenAvailableを使用するには、そのデータベース内のすべてのコレクションに対して変更ストリームの事前イメージと事後イメージを有効にする必要があります。

config.fullDocumentOnly

ブール値

条件付き

すべてのメタデータを含む変更イベント ドキュメント全体を返すか、 fullDocumentの内容のみを返すかを制御する設定。 trueに設定されている場合、ソースはfullDocumentの内容のみを返します。

がfullDocument requiredまたはwhenAvailable に設定されている場合は、そのデータベース内のすべてのコレクションに対して変更ストリームの事前イメージと事後イメージを有効にする必要があります。

config.fullDocumentBeforeChange

string

任意

変更ストリーム ソースに、出力に元の「変更前」状態の完全なドキュメントを含めるかどうかを指定します。 次のいずれかである必要があります。

  • off : fullDocumentBeforeChangeフィールドを省略します。

  • required : 状態が変更される前に、完全なドキュメントを返す必要があります。 状態が変更される前の完全なドキュメントが利用できない場合、ストリーム プロセッサは失敗します。

  • whenAvailable : 使用可能な場合は常に、変更前の状態で完全なドキュメントを返します。それ以外の場合は、 fullDocumentBeforeChangeフィールドを省略します。

fullDocumentBeforeChangeの値を指定しない場合、デフォルトはoffになります。

データベース変更ストリームでこのフィールドを使用するには、そのデータベース内のすべてのコレクションに対して変更ストリームの事前イメージと事後イメージを有効にする必要があります。

config.pipeline

ドキュメント

任意

変更ストリーム出力を発生元でフィルタリングするための集計パイプラインを指定します。このパイプラインは「変更ストリーム出力の修正」で説明されているパラメータに準拠する必要があります。

重要: 各変更イベントには wallTime フィールドと clusterTime フィールドが含まれます。$source 以降の Atlas Stream Processing ステージでは、プロセッサがこれらのフィールドを取り込んだため、受信することが想定されています。Change Stream データの適切なプロセシングを確保するために、$source.config.pipeline ではこれらのフィールドを変更しないでください。

config.maxAwaitTimeMS

integer

任意

空のバッチを返す前に、新しいデータ変更が変更ストリーム カーソルに報告されるまで待機する最大時間(ミリ秒)。

デフォルトは 1000 です。

変更ストリーム$sourceを使用する場合は、ソースクラスターで少なくとも24時間のoplog windowを構成してください。

変更ストリームを読み取るために、Atlas Stream Processing は oplog コレクションをスキャンします。その結果、ログに COLLSCAN 警告が表示されることがあります。これらの警告は通常の動作を示しており、エラーを示すものではありません。

config.fullDocumentまたはconfig.fullDocumentBeforeChangeをrequiredに設定する場合は、キャプチャー対象の書き込み操作を実行する前に、各コレクションでchangeStreamPreAndPostImagesを有効にしてください。書き込み発生時に機能が有効化されていなかった場合、またはpost-imageの有効期限が切れていた場合に、イベントでpost-imageを利用できないと、ストリームプロセッサーは失敗します。変更前ドキュメントおよび変更後ドキュメントを有効化する方法については、「ドキュメントのPre-ImageおよびPost-Imageを使用したChange Streams」を参照してください。

重要

変更ストリームを再開できるのは、再開トークンが識別する操作をoplogコレクションが保持している間のみです。アプリケーションが変更ストリームを再開できるようにするには、予想される最長時間よりも長い最小 Oplog Window を設定します。

Atlas クラスター変更ストリーム全体からのストリーミングデータを操作するには、$source ステージのプロトタイプ形式は次のようになります。

{
"$source": {
"connectionName": "<registered-connection>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"readPreference": "<read-preference>",
"readPreferenceTags": [
{"<key>": "<value>"},
. . .
],
"config": {
"startAfter": <start-token> | "startAtOperationTime": <timestamp>,
"fullDocument": "<full-doc-condition>",
"fullDocumentOnly": <boolean>,
"fullDocumentBeforeChange": "<before-change-condition>",
"pipeline": [{
"<aggregation-stage>" : {
<stage-input>,
. . .
},
. . .
}]
},
}
}

$sourceステージは、次のフィールドを持つドキュメントを取得します。

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

connectionName

string

条件付き

データを取り込む接続レジストリ内の接続を識別するラベル。

timeField

ドキュメント

任意

受信メッセージの権限のあるタイムスタンプを定義するドキュメント。

timeFieldを使用する場合は、次のいずれかとして定義する必要があります。

  • ソース メッセージ フィールドを引数として$toDate式

  • ソース メッセージ フィールドを引数として$dateFromString式。

timeFieldを宣言しない場合、Atlas Stream Processing は、ソースによって提供されたメッセージ タイムスタンプからタイムスタンプを作成します。

readPreference

string

任意

readPreferenceTags

配列

任意

config

ドキュメント

任意

のさまざまなデフォルト値を上書きするフィールドを含むドキュメント。

config.startAfter

token

条件付き

ソースがレポートを開始する変更イベント。 これは再開トークンの形式をとります。

config.startAfterまたはconfig.startAtOperationTimeのいずれか 1 つだけを使用できます。

config.startAtOperationTime

日付 | タイムスタンプ

条件付き

ソースがレポートを開始するoptime 。

config.startAfterまたはconfig.startAtOperationTimeのいずれか 1 つだけを使用できます。

MongoDB 拡張 JSON の $date または $timestamp の値を受け付けます。

config.fullDocument

string

条件付き

変更ストリーム ソースが完全なドキュメントを返すか、更新が発生したときにのみ変更を返すかを制御する設定。 次のいずれかである必要があります。

  • default : サーバーのデフォルトの動作を使用します。update 操作の場合、完全なドキュメントは返されません。

  • updateLookup : 更新操作によって行われた変更に加えて、過半数がコミットした現在のバージョンの完全なドキュメントを返します。

  • required : 完全なドキュメントを返す必要があります。完全なドキュメントが利用できない場合、ストリームプロセッサは失敗します。これは削除操作には適用されません。この操作では、完全なドキュメントが利用できない場合に削除イベントが作成されます。

  • whenAvailable : 完全なドキュメントが利用可能になるたびに完全なドキュメントを返し、そうでない場合は変更を返します。

fullDocumentの値を指定しない場合、デフォルトはdefaultになります。

クラスター変更ストリームで またはrequired whenAvailableを使用するには、そのクラスター内のすべてのコレクションに対して変更ストリームの事前イメージと事後イメージを有効にする必要があります。

config.fullDocumentOnly

ブール値

条件付き

すべてのメタデータを含む変更イベント ドキュメント全体を返すか、 fullDocumentの内容のみを返すかを制御する設定。 trueに設定されている場合、ソースはfullDocumentの内容のみを返します。

がfullDocument requiredまたはwhenAvailable に設定されている場合、そのクラスター内のすべてのコレクションで変更ストリームの事前イメージと事後イメージを有効にする必要があります。

config.fullDocumentBeforeChange

string

任意

変更ストリーム ソースに、出力に元の「変更前」状態の完全なドキュメントを含めるかどうかを指定します。 次のいずれかである必要があります。

  • off : fullDocumentBeforeChangeフィールドを省略します。

  • required : 状態が変更される前に、完全なドキュメントを返す必要があります。 状態が変更される前の完全なドキュメントが利用できない場合、ストリーム プロセッサは失敗します。

  • whenAvailable : 使用可能な場合は常に、変更前の状態で完全なドキュメントを返します。それ以外の場合は、 fullDocumentBeforeChangeフィールドを省略します。

fullDocumentBeforeChangeの値を指定しない場合、デフォルトはoffになります。

データベース変更ストリームでこのフィールドを使用するには、そのデータベース内のすべてのコレクションに対して変更ストリームの事前イメージと事後イメージを有効にする必要があります。

config.pipeline

ドキュメント

任意

変更ストリーム出力を発生元でフィルタリングするための集計パイプラインを指定します。このパイプラインは「変更ストリーム出力の修正」で説明されているパラメータに準拠する必要があります。

Atlas Stream Processing は、取り込まれた各変更イベントから wallTime フィールドと clusterTime フィールドを受信することが予想されていることに注意してください。Change Stream データの適切な処理を確保するために、$source.config.pipeline ではこれらのフィールドを変更しないでください。

config.maxAwaitTimeMS

integer

任意

空のバッチを返す前に、新しいデータ変更が変更ストリーム カーソルに報告されるまで待機する最大時間(ミリ秒)。

デフォルトは 1000 です。

Atlas Stream Processing は、 AWS Kinesisストリームへの Private Link 接続 の作成をサポートしています。詳細については、Kinesis Private Link 接続の追加 を参照してください。

AWS Kinesisデータストリームのデータを操作する場合、$source ステージには次のプロトタイプ形式があります。

{
"$source": {
"connectionName": "<registered-connection>",
"stream": "<stream-name>",
"region": "<aws-region>",
"timeField": {
$toDate | $dateFromString: <expression>
},
"tsFieldName": "<field-name>",
"shardIdleTimeout": {
"size": <duration-number>,
"unit": "<duration-unit>"
},
"config": {
"consumerARN": "<aws-arn>",
"initialPosition": <initial-position>,
reshardDetectionIntervalSecs: <interval>
}
}
}

$sourceステージは、次のフィールドを持つドキュメントを取得します。

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

connectionName

string

必須

データを取り込む接続レジストリ内の接続を識別するラベル。

config.consumerARN

string

必須

stream

string

必須

メッセージをストリーミングするAWS Kinesisデータストリーム。

region

string

条件付き

指定されたストリームが存在するAWSリージョン。Kinesis は、異なるリージョンで同じ名前の複数のデータ ストリームをサポートしています。同じ接続内の 2 つ以上のリージョンのデータストリームに同じ名前を使用する場合は、使用する名前とリージョンの組み合わせを指定するために、このフィールドを使用する必要があります。

timeField

ドキュメント

任意

受信メッセージの権限のあるタイムスタンプを定義するドキュメント。

timeFieldを使用する場合は、次のいずれかとして定義する必要があります。

  • ソース メッセージ フィールドを引数として$toDate式

  • ソース メッセージ フィールドを引数として$dateFromString式。

timeFieldを宣言しない場合、Atlas Stream Processing は、ソースによって提供されたメッセージ タイムスタンプからタイムスタンプを作成します。

tsFieldName

string

任意

プロジェクションされたドキュメントのタイムスタンプのフィールド名。このフィールドを使用して、デフォルトのタイムスタンプフィールド名を上書きします。

shardIdleTimeout

ドキュメント

任意

埋め込みの計算で無視される前に、シャードがアイドル状態になる時間を指定するドキュメント。

このフィールドは、デフォルトで無効になっています。アイドル状態であるため、前に移動しないシャードを取り扱うには、このフィールドに値を設定します。

shardIdleTimeout.size

ドキュメント

任意

シャード アイドル タイムアウトの期間を指定する数値。

shardIdleTimeout.unit

ドキュメント

任意

シャード アイドル タイムアウトの期間の単位。

unitの値は次のいずれかになります。

  • "ms" (ミリ秒)

  • "second"

  • "minute"

  • "hour"

  • "day"

config

ドキュメント

任意

のさまざまなデフォルト値を上書きするフィールドを含むドキュメント。

config.initialPosition

string

任意

Kinesis データ ストリームの履歴内の、メッセージの取り込みを開始する位置。次のいずれかである必要があります。

  • "TRIM_HORIZON": シャード内の最も古いメッセージから取り込みを開始します。

  • "LATEST": シャード内の最新のメッセージから取り込みを開始します。

  • "AT_TIMESTAMP": 特定のタイムスタンプから取り込みを開始します。config.atTimestampが必要です。

  • "AT_SEQUENCE_NUMBER": 特定のシーケンス番号で取り込みを開始します。

  • "AFTER_SEQUENCE_NUMBER": 特定のシーケンス番号の後に取り込みを開始します。

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

config.atTimestamp

date

条件付き

メッセージの取り込みを開始するタイムスタンプ。config.initialPosition が "AT_TIMESTAMP" の場合は必須です。

config.reshardDetectionIntervalSecs

integer

任意

ドキュメントの配列を操作するために、 $sourceステージには次のプロトタイプ形式があります。

{
"$source": {
"timeField": {
$toDate | $dateFromString: <expression>
},
"documents" : [{source-doc},...] | <expression>
}
}

$sourceステージは、次のフィールドを持つドキュメントを取得します。

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

timeField

ドキュメント

任意

受信メッセージの権限のあるタイムスタンプを定義するドキュメント。

timeFieldを使用する場合は、次のいずれかとして定義する必要があります。

  • ソース メッセージ フィールドを引数として$toDate式

  • ソース メッセージ フィールドを引数として$dateFromString式。

timeFieldを宣言しない場合、Atlas Stream Processing は、ソースによって提供されたメッセージ タイムスタンプからタイムスタンプを作成します。

documents

配列

条件付き

ストリーミング データソースとして使用するドキュメントの配列。 このフィールドの値は、オブジェクトの配列、またはオブジェクトの配列として評価される 式 のいずれかになります。 connectionNameフィールドを使用する場合は、このフィールドを使用しないでください。

接続からの読み取りではなく、定期的なスケジュールでドキュメントを生成する場合、$source ステージのプロトタイプ形式は次のとおりです。

{
"$source": {
"schedule": "<cron-expression>",
"tsFieldName": "<timestamp-field-name>"
}
}

$sourceステージは、次のフィールドを持つドキュメントを取得します。

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

schedule

string

必須

ステージがドキュメントをいつ発行するかを決定する 6 フィールドの CR式。フィールドは、秒、分、時間、日付、月、曜日の順に並べられます。 Atlas Stream Processing は UTC で式を評価します。

各フィールドは、すべての値、単一の値、1-5 などの範囲、0/15 などのステップ、または 1,15,30 などのカンマ区切りリストに対して * を受け入れます。月には JANから DEC まで、 曜日には SUN ~ SAT という名前を使用できます。 CRON_TZ プレフィックスを含めないでください。

Atlas Stream Processing が、日付と曜日の両方を制限する式を解決する方法の詳細については、「 動作 」を参照してください。

tsFieldName

string

任意

ステージがスケジュールされたタイムスタンプをプロジェクションするフィールドの名前。デフォルトは _ts です。

$sourceは、表示されるすべてのパイプラインの最初のステージである必要があります。 パイプラインごとに使用できる$sourceステージは 1 つだけです。

Kafka $source ステージでは、Atlas Stream Processing はソーストピック内の複数のパーティションから並列に読み取りを行います。パーティションの制限は、プロセッサ階層によって決まります。詳しくは、Stream Processing 請求参照を参照してください。

cron $source ステージでは、スケジュールされた時間ごとに、スケジュールされたタイムスタンプを保持する空のドキュメントが 1 つ生成されます。パイプラインの後半のステージを使用して、ドキュメント にドキュメントを入力します。

schedule式が日付と曜日の両方を制限する場合、Atlas Stream Processing は両方のフィールドを満たすスケジュールされた時刻にのみドキュメントを発行します。これは、いずれかのフィールドが一致したときにドキュメントを発行する cron 実装とは異なります。

cron $source ステージのストリーム プロセッサが再起動すると、 を実行中いない間に失敗したスケジュールされた時間ごとにドキュメントが発行されます。停止後に欠落したスケジュールされた時間のドキュメントは発行されません。

cron$source ステージを持つストリーム プロセッサはスケジュールされた時間間で継続的に実行され、Atlas はその期間にわたってプロセッサ階層にそれを請求します。頻度の低いスケジュールでは、プロセッサの実行中コストは削減されません。詳しくは、「 Atlas Stream Processing の請求参照 」を参照してください。

ストリーミング データソースは、気象用サンプル データセットのスキーマに準拠して、さまざまなロケーションから詳細な気象レポートを生成します。次の集計には 3 つのステージがあります。

  1. $sourceステージでは、 という名前のトピックでこれらのレポートを収集するApachemy_weatherdata Kafkaプロバイダーとの接続を確立し、各レコードが後続の集計ステージに取り込まれる際に公開します。このステージではまた、プロジェクションを実行するタイムスタンプフィールドの名前が上書きされ、ingestionTime に設定されます。

  2. $matchステージでは、dewPoint.value 5.0が 以下のドキュメントを除外し、 がdewPoint.value 5.0より大きいドキュメントを次のステージに渡します。

  3. $mergeステージは、 sample_weatherstreamデータベース内のstreamという名前の Atlas コレクションに出力を書き込みます。 そのようなデータベースやコレクションが存在しない場合は、Atlas がそれらを作成します。

[{
"$source": {
"connectionName": "sample_weatherdata",
"topic": "my_weatherdata"
}
},
{
"$match": { "dewPoint.value": { "$gt": 5 } }
},
{
"$merge": {
"into": {
"connectionName": "weatherStreamOutput",
"db": "sample_weatherstream",
"coll": "stream"
}
}
}]

結果のsample_weatherstream.streamコレクション内のドキュメントを表示するには、Atlas クラスターに接続して次のコマンドを実行します。

db.getSiblingDB("sample_weatherstream").stream.find()

注意

前述の例はその一般的な例です。 ストリーミング データは静的ではなく、各ユーザーに異なるドキュメントが表示されます。

次の集計は、cluster0-collection ソースからデータを取り込み、サンプルデータセットがロードされた Atlas クラスターに接続します。Stream Processing ワークスペースを作成し、Atlas クラスターへの接続を接続レジストリに追加する方法については、Atlas Stream Processing の始め方 をご覧ください。この集計は 2 つのステージを実行して、sample_weatherdata データベースの data コレクションに対する変更ストリームを開き、変更を記録します。

  1. $source ステージは cluster0-collection ソースに接続し、sample_weatherdata データベース内の data コレクションに対して変更ストリームを開きます。

  2. $merge ステージは、フィルタリングされた変更ストリームドキュメントを、sample_weatherdata データベース内の data_changes という名前の Atlas コレクションに書き込みます。そのようなコレクションが存在しない場合、Atlas が作成します。

[{
"$source": {
"connectionName": "cluster0-connection",
"db": "sample_weatherdata",
"coll": "data"
}
},
{
"$merge": {
"into": {
"connectionName": "cluster0-connection",
"db": "sample_weatherdata",
"coll": "data_changes"
}
}
}]

次の mongosh コマンドは data ドキュメントを削除します。

db.getSiblingDB("sample_weatherdata").data.deleteOne(
{ _id: ObjectId("5553a99ae4b02cf715120e4b") }
)

data ドキュメントが削除された後、ストリームプロセッサは変更ストリームイベントドキュメントを sample_weatherdata.data_changes コレクションに書き込みます。結果の sample_weatherdata.data_changes コレクション内のドキュメントを表示するには、mongosh を使用して Atlas クラスターに接続し、次のコマンドを実行してください。

db.getSiblingDB("sample_weatherdata").data_changes.find()

次の集計では、サンプルdb-change-stream-connection Mflix データセット コレクション データセットがロードされた Atlas クラスターに接続する ソースからデータを取り込みます。ストリーム処理ワークスペースを作成し、Atlas クラスターへの接続を接続レジストリに追加する方法については、「 Atlas Stream Processing を使い始める 」を参照してください。この集計では2 つのステージを実行して、sample_mflix ソースデータベースに対して変更ストリームを開き、 シンクデータベースのdb_changes sample_mflix_changesコレクションへの変更をレコード。

  1. $sourceステージはdb-change-stream-connection ソースに接続し、sample_mflix ソースデータベースに対して変更ストリームを開きます。config.startAtOperationTime フィールドは、ソースがレポートを開始する時間を設定します。この例では、startAt の値を 1 分前に開始するように設定します。

  2. $mergeステージは、変更ストリームドキュメントをdb_changes sample_mflix_changesSink データベース内の という名前の Atlasコレクションに書き込みます。

const startAt = new Date(Date.now() - 60 * 1000);
const pipeline = [
{
$source: {
connectionName: "db-change-stream-connection",
db: "sample_mflix",
config: {
startAtOperationTime: startAt
}
}
},
{
$merge: {
into: {
connectionName: "db-change-stream-connection",
db: "sample_mflix_changes",
coll: "db_changes"
}
}
}
];

ソースクラスターに対して次のコマンドを実行し、ストリームプロセッサの動作を確認します。接続するには、 「 mongosh経由でクラスターに接続する 」を参照してください。

sample_mflix ソースデータベース内の moviesコレクションにドキュメントを挿入し、commentsコレクションにドキュメントを挿入します。

db.getSiblingDB("sample_mflix").movies.insertOne({
title: "The Stream Processor",
year: 2026
})
db.getSiblingDB("sample_mflix").comments.insertOne({
name: "Ada Lovelace",
text: "A fine film about data in motion."
})

ドキュメントを挿入すると、ストリーム プロセッサは挿入ごとに変更ストリームイベントドキュメントをsample_mflix_changes.db_changesコレクションに書込みます。クラスターに対して次のコマンドを実行すると、結果として得られる sample_mflix_changes.db_changesコレクション内のドキュメントを表示します。

db.getSiblingDB("sample_mflix_changes").db_changes.find(
{},
{ _id: 0, clusterTime: 1, ns: 1, operationType: 1, fullDocument: 1 }
)

各イベントには、タイムスタンプ、ドキュメント全体の変更、operationType、nsフィールドが含まれます。 nsフィールドはソースデータベースとコレクション を名前付きするため、各変更を生成したコレクションを区別できます。

次の集計では、サンプルcluster-changestream-connection Mflix Dataset Collections データセットにロードされたcluster-changestream-sink-connection Atlas クラスターに接続する ソースからデータを取り込み、別の Atlas クラスターに接続する 宛先に書き込みます。ストリーム処理ワークスペースを作成し、Atlas クラスターへの接続を接続レジストリに追加する方法については、「 Atlas Stream Processing を使い始める 」を参照してください。この集計では 2 つのステージを実行して、ソースクラスターでクラスター全体の変更ストリームを開き、宛先クラスターのevents cluster_changesデータベースにある コレクションへの変更をレコード。この$source ステージではdb とcoll が省略されているため、単一のコレクションやデータベースではなく、ソースクラスター上のすべてのデータベースとコレクションからの変更が報告されます。

  1. $sourcecluster-changestream-connectionステージは ソースに接続し、ソースクラスター全体に対して変更ストリームを開きます。 フィールドは、ストリームconfig.startAtOperationTime プロセッサが または 指定した時間の後に発生する変更のレポートを開始することを指定します。

  2. $mergeステージでは、宛先クラスターのevents cluster_changesデータベース内の という名前の Atlasコレクションに変更ストリームドキュメントが書き込まれます。そのようなデータベースやコレクションが存在しない場合は、Atlas がそれらを作成します。

[{
"$source": {
"connectionName": "cluster-changestream-connection",
"config": {
"startAtOperationTime": {"$date": "2024-08-19T18:00:00.000Z"}
}
}
},
{
"$merge": {
"into": {
"connectionName": "cluster-changestream-sink-connection",
"db": "cluster_changes",
"coll": "events"
}
}
}]

mongosh次のmovies sample_mflixコマンドは、ソースクラスター上の データベース内の コレクションにドキュメントを挿入します。

db.getSiblingDB("sample_mflix").movies.insertOne(
{ _id: ObjectId("66c1a1f1f1f1f1f1f1f1f1f1"), title: "Example Movie" }
)

ドキュメントが挿入されると、ストリーム プロセッサは変更ストリームイベントドキュメントを宛先クラスターの cluster_changes.eventsコレクションに書込みます。

結果のcluster_changes.events コレクション内のドキュメントを表示するには、宛先クラスターに対して次のコマンドを実行します。接続するには、 「 mongosh経由でクラスターに接続する 」を参照してください。

db.getSiblingDB("cluster_changes").events.find()

出力ドキュメントの nsフィールドは、変更がソースクラスターの sample_mflix.movies 挿入からのものであることを示します。変更ストリーム集計パイプラインは、宛先クラスターの cluster_changes.eventsコレクションへのこの変更を反映します。

次の集計では、インラインドキュメント配列をストリーミングデータソースとして使用し、3 つのロケーションの気象観測が含まれています。配列は、気象用サンプル データセットと同じスキーマを使用します。この集計は 3 つのステージを実行します。

  1. $source ステージでは、気象測定のインライン documents 配列をストリーミングデータソースとして定義し、timeField を使用して各ドキュメントの timestampフィールドを認証タイムスタンプとして指定します。

  2. $match ステージでは、dewPoint.value が 5.0 を超えるドキュメントのみが次のステージに渡されます。

  3. $mergeステージは、 sample_weatherstreamデータベース内のstreamという名前の Atlas コレクションに出力を書き込みます。 そのようなデータベースやコレクションが存在しない場合は、Atlas がそれらを作成します。

[{
"$source": {
"documents": [
{
"location": "New York",
"timestamp": ISODate('2024-01-15T08:00:00Z'),
"temp": 23.5,
"dewPoint": { "value": 6.2 }
},
{
"location": "Los Angeles",
"timestamp": ISODate('2024-01-15T08:05:00Z'),
"temp": 18.2,
"dewPoint": { "value": 4.8 }
},
{
"location": "Chicago",
"timestamp": ISODate('2024-01-15T08:10:00Z'),
"temp": 26.8,
"dewPoint": { "value": 7.5 }
}
]
}
},
{
"$match": { "dewPoint.value": { "$gt": 5.0 } }
},
{
"$merge": {
"into": {
"connectionName": "weatherStreamOutput",
"db": "sample_weatherstream",
"coll": "stream"
}
}
}]

結果のsample_weatherstream.streamコレクション内のドキュメントを表示するには、Atlas クラスターに接続して次のコマンドを実行します。

db.getSiblingDB("sample_weatherstream").stream.find()

sample_stream_solar次の集計フィルターは、サンプルストリーミングデータソース からのレポートを出力します。結果はsolar-cluster-connection 接続にアーカイブされます。ストリーム処理ワークスペースを作成し、Atlas Stream Processing ワークスペースに接続を追加する方法については、「 Atlas Stream Processing の使用を開始 」と「 Atlas Stream Processing 接続の追加 」を参照してください。この集計では3 つのステージを実行して、sample_stream_solar ソースからのレポートをフィルタリングし、その結果をsolarColl solarDbデータベース内の という名前のコレクションに書き込みます。

  1. $sourcesample_stream_solarステージは ソースに接続します。 フィールドは、timeField timestampを使用して各受信レポートの$dateFromString フィールドを日付に変換します。

  2. $matchステージでは、device_id がdevice_8 であるドキュメントを除外し、他のすべてのデバイスからのレポートを次のステージに渡します。

  3. ステージは、 $mergeクラスターの データベース内のsolarColl コレクションに出力を書き込みます。solarDbsolar-cluster-connection

[{
"$source": {
"connectionName": "sample_stream_solar",
"timeField": {
"$dateFromString": { "dateString": "$timestamp" }
}
}
},
{
"$match": { "device_id": { "$ne": "device_8" } }
},
{
"$merge": {
"into": {
"connectionName": "solar-cluster-connection",
"db": "solarDb",
"coll": "solarColl"
}
}
}]

注意

sample_stream_solar ソースは、1 秒ごとにサンプルドキュメントを生成し、高速プロトタイプ作成を目的としたテスト専用の接続です。

solar-cluster-connectionクラスターに対して次のコマンドを実行して、ストリーム プロセッサの動作を確認します。接続するには、 「 mongosh経由でクラスターに接続する 」を参照してください。

結果の solarDb.solarCollコレクション内のドキュメントを表示するには、次のコマンドを実行します。

db.getSiblingDB("solarDb").solarColl.find()

各ドキュメントには、ソート デバイス レポートのdevice_id event_type、 、group_id 、max_watts 、obs 、timestamp フィールドが含まれています。出力内のドキュメントにも が でないことがわかります。これは、device_id device_8$matchステージが$merge ステージが残りのドキュメントをsolarColl に書き込む前にそれらのレポートを除外するためです。

注意

前述の例はその一般的な例です。 ストリーミング データは静的ではなく、各ユーザーに異なるドキュメントが表示されます。

次の集計では、接続からの読み取りではなく、5 分ごとにドキュメントが生成されます。この集計は3 つのステージを実行します。

  1. $sourceステージでは0 5 分ごとの秒 に空のドキュメントが生成されます(UTC)。

  2. $projectステージでは、各ドキュメントにジョブ名がラベル付けされ、_ts フィールドから フィールドにスケジュールされたタイムスタンプがコピーされます。runAt

  3. $mergeステージは、 sample_weatherstreamデータベース内のheartbeatsという名前の Atlas コレクションに出力を書き込みます。 そのようなデータベースやコレクションが存在しない場合は、Atlas がそれらを作成します。

[{
"$source": {
"schedule": "0 0/5 * * * *"
}
},
{
"$project": {
"job": "five-minute-heartbeat",
"runAt": "$_ts"
}
},
{
"$merge": {
"into": {
"connectionName": "weatherStreamOutput",
"db": "sample_weatherstream",
"coll": "heartbeats"
}
}
}]

結果のsample_weatherstream.heartbeatsコレクション内のドキュメントを表示するには、Atlas クラスターに接続して次のコマンドを実行します。

db.getSiblingDB("sample_weatherstream").heartbeats.find()

注意

前の は の の例です。表示されるタイムスタンプは、ストリーム プロセッサを起動するときに異なります。