定義
$sourceステージでは、データをストリーミングするための接続を接続レジストリで指定します。 次の接続タイプがサポートされています。
Apache Kafkaエージェント
MongoDB コレクションの変更ストリーム
MongoDB database 変更ストリーム
MongoDBクラスターの変更ストリーム
AWS Kinesisデータストリーム
ドキュメント配列
CRON スケジュール
構文
Apache Kafka ブロック
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ステージは、次のフィールドを持つドキュメントを取得します。
フィールド | タイプ | 必要性 | 説明 | |
|---|---|---|---|---|
| string | 必須 | データを取り込む接続レジストリ内の接続を識別するラベル。 | |
| 文字列または複数の文字列の配列 | 必須 | メッセージをストリーミングする 1 つ以上の Apache Kafka トピックの名前。複数のトピックからのメッセージをストリーミングする場合は、配列で指定します。 | |
| ドキュメント | 任意 | 受信メッセージの権限のあるタイムスタンプを定義するドキュメント。
| |
| ドキュメント | 任意 | 証明機関の計算で無視される前に、パーティションがアイドル状態になることを許可する時間を指定するドキュメント。 このフィールドはデフォルトで無効になっています。アイドル状態で進まないパーティションを処理するには、このフィールドに値を設定します。 | |
| integer | 任意 | パーティションのアイドル タイムアウトの期間を指定する数値。 | |
| string | 任意 | パーティション アイドル タイムアウトの期間の単位。
| |
| ドキュメント | 任意 | 平均直列化されたソースからの読み取りをサポートするためにスキーマ レジストリの使用を有効にするドキュメント。 この機能を有効にするには、スキーマ Registry 接続を作成する必要があります。 | |
| string | 条件付き | Avro 非直列化に使用するスキーマ レジストリ接続の名前。 | |
| ドキュメント | 任意 | のさまざまなデフォルト値を上書きするフィールドを含むドキュメント。 | |
| string | 任意 | Apache Kafkaソーストピック内のどのイベントで取り込みを開始するかを指定します。
デフォルトは | |
| string | 任意 | ストリーム プロセッサと関連付ける Kafka コンシューマー グループの ID。省略した場合、Atlas Stream Processing は、Stream Processing ワークスペースを次の形式の自動生成された ID に関連付けます。 Atlas Stream Processing は、すべての永続ストリーム プロセッサに対してこのパラメーターの値を自動的に生成します。SP.process() で定義されたエフェメラル ストリーム プロセッサの場合、このパラメーターは手動で定義した場合にのみ設定されます。 | |
| ブール値 | 条件付き | Kafkaプロバイダー パーティション オフセットのコミット ポリシーを決定するフラグ。Atlas Stream Processing は 2 つのコミット ポリシーをサポートしています。
SP.process() で定義されたエフェメラル ストリーム プロセッサの場合、 Kafka を | |
| string | 任意 | Apache Kafkaキー データを逆直列化するために使用されるデータ型。次のいずれかの値である必要があります。
デフォルトは | |
| string | 任意 | Apache Kafkaキー データを逆直列化するときに発生したエラーの処理方法。次のいずれかの値である必要があります。
|
注意
Atlas Stream Processing では、ソース データ ストリーム内のドキュメントが有効なjsonまたはejsonである必要があります。 Atlas Stream Processing は、この要件を満たさないドキュメントをデッド レター キューに設定します(デッド レター キューを設定している場合)。
MongoDB コレクションの変更ストリーム
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ステージは、次のフィールドを持つドキュメントを取得します。
フィールド | タイプ | 必要性 | 説明 |
|---|---|---|---|
| string | 条件付き | データを取り込む接続レジストリ内の接続を識別するラベル。 |
| ドキュメント | 任意 | 受信メッセージの権限のあるタイムスタンプを定義するドキュメント。
|
| string | 必須 |
|
| 文字列または複数の文字列の配列 | 必須 | によって指定された Atlasインスタンスでホストされている 1 つ以上のMongoDBコレクションの名前。これらのコレクションの変更ストリームは、ストリーミングデータソースとして機能します。このフィールドを省略すると、ストリーム |
| ドキュメント | 任意 |
Atlas Stream Processing
Atlas Stream Processing は、ストリーム プロセッサの起動時に同期するコレクションのセットを修正します。点の後に作成するコレクションは同期されません。いずれかのコレクションの同期に失敗すると、ストリーム プロセッサ全体が失敗します。
|
| ブール値 | 条件付き |
|
| integer | 任意 |
設定できる最大値は、ストリームプロセッサーの階層によって異なります。詳しくは、Atlas Stream Processing 階層選択ガイドを参照してください。 各ストリーム プロセッサには、その階層によって決定される最大累積並列処理値があります。ストリーム プロセッサの累積並列処理は、次のように計算されます。
例、 ストリーム プロセッサがその階層の最大累積並列処理を超える場合、Atlas Stream Processing はエラーをスローし、目的のレベルの並列処理に必要な最小プロセッサ階層について提案します。エラーを解決するには、プロセッサをより高い階層に増やすアップするか、ステージの並列処理値を低くする必要があります。詳しくは、Stream Processingをご覧ください。 |
| string | 任意 | 読み込み設定 (read preference)変更ストリームと デフォルトは |
| 配列 | 任意 | 読み込み設定 (read preference) タグ変更ストリームと |
| ドキュメント | 任意 | のさまざまなデフォルト値を上書きするフィールドを含むドキュメント。 |
| token | 条件付き | ソースがレポートを開始する変更イベント。 これは再開トークンの形式をとります。
|
| タイムスタンプ | date | 条件付き | ソースがレポートを開始するoptime 。
MongoDB 拡張 JSON の |
| string | 条件付き | 変更ストリーム ソースが完全なドキュメントを返すか、更新が発生したときにのみ変更を返すかを制御する設定。 次のいずれかである必要があります。
コレクション変更ストリームで または |
| ブール値 | 条件付き | |
| string | 任意 | 変更ストリーム ソースに、出力に元の「変更前」状態の完全なドキュメントを含めるかどうかを指定します。 次のいずれかである必要があります。
コレクションの変更ストリームでこのフィールドを使用するには、そのコレクションで変更ストリームの事前イメージと事後イメージを有効にする必要があります。 |
| ドキュメント | 任意 | 集計パイプラインを指定し、変更ストリーム出力がパスされ、さらなるプロセシングが実行される前にフィルタリングします。このパイプラインは「変更ストリーム出力の修正」で説明されているパラメータに準拠する必要があります。 重要: 各変更イベントには |
| integer | 任意 | 空のバッチを返す前に、新しいデータ変更が変更ストリーム カーソルに報告されるまで待機する最大時間(ミリ秒)。 デフォルトは |
MongoDB Database Change Stream
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ステージは、次のフィールドを持つドキュメントを取得します。
フィールド | タイプ | 必要性 | 説明 |
|---|---|---|---|
| string | 条件付き | データを取り込む接続レジストリ内の接続を識別するラベル。 |
| ドキュメント | 任意 | 受信メッセージの権限のあるタイムスタンプを定義するドキュメント。
|
| string | 必須 |
|
| string | 任意 | 変更ストリーム操作 の読み込み設定(read preference)。 デフォルトは |
| 配列 | 任意 | |
| ドキュメント | 任意 | のさまざまなデフォルト値を上書きするフィールドを含むドキュメント。 |
| token | 条件付き | ソースがレポートを開始する変更イベント。 これは再開トークンの形式をとります。
|
| タイムスタンプ | date | 条件付き | ソースがレポートを開始するoptime 。
MongoDB 拡張 JSON の |
| string | 条件付き | 変更ストリーム ソースが完全なドキュメントを返すか、更新が発生したときにのみ変更を返すかを制御する設定。 次のいずれかである必要があります。
データベース変更ストリームで または |
| ブール値 | 条件付き | |
| string | 任意 | 変更ストリーム ソースに、出力に元の「変更前」状態の完全なドキュメントを含めるかどうかを指定します。 次のいずれかである必要があります。
データベース変更ストリームでこのフィールドを使用するには、そのデータベース内のすべてのコレクションに対して変更ストリームの事前イメージと事後イメージを有効にする必要があります。 |
| ドキュメント | 任意 | 変更ストリーム出力を発生元でフィルタリングするための集計パイプラインを指定します。このパイプラインは「変更ストリーム出力の修正」で説明されているパラメータに準拠する必要があります。 重要: 各変更イベントには |
| integer | 任意 | 空のバッチを返す前に、新しいデータ変更が変更ストリーム カーソルに報告されるまで待機する最大時間(ミリ秒)。 デフォルトは |
MongoDB クラスター全体の変更ストリームソース
変更ストリーム$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ステージは、次のフィールドを持つドキュメントを取得します。
フィールド | タイプ | 必要性 | 説明 |
|---|---|---|---|
| string | 条件付き | データを取り込む接続レジストリ内の接続を識別するラベル。 |
| ドキュメント | 任意 | 受信メッセージの権限のあるタイムスタンプを定義するドキュメント。
|
| string | 任意 | 変更ストリーム操作 の読み込み設定(read preference)。 デフォルトは |
| 配列 | 任意 | |
| ドキュメント | 任意 | のさまざまなデフォルト値を上書きするフィールドを含むドキュメント。 |
| token | 条件付き | ソースがレポートを開始する変更イベント。 これは再開トークンの形式をとります。
|
| 日付 | タイムスタンプ | 条件付き | ソースがレポートを開始するoptime 。
MongoDB 拡張 JSON の |
| string | 条件付き | 変更ストリーム ソースが完全なドキュメントを返すか、更新が発生したときにのみ変更を返すかを制御する設定。 次のいずれかである必要があります。
クラスター変更ストリームで または |
| ブール値 | 条件付き | |
| string | 任意 | 変更ストリーム ソースに、出力に元の「変更前」状態の完全なドキュメントを含めるかどうかを指定します。 次のいずれかである必要があります。
データベース変更ストリームでこのフィールドを使用するには、そのデータベース内のすべてのコレクションに対して変更ストリームの事前イメージと事後イメージを有効にする必要があります。 |
| ドキュメント | 任意 | 変更ストリーム出力を発生元でフィルタリングするための集計パイプラインを指定します。このパイプラインは「変更ストリーム出力の修正」で説明されているパラメータに準拠する必要があります。 Atlas Stream Processing は、取り込まれた各変更イベントから |
| integer | 任意 | 空のバッチを返す前に、新しいデータ変更が変更ストリーム カーソルに報告されるまで待機する最大時間(ミリ秒)。 デフォルトは |
AWS Kinesis Data Stream
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ステージは、次のフィールドを持つドキュメントを取得します。
フィールド | タイプ | 必要性 | 説明 |
|---|---|---|---|
| string | 必須 | |
| string | 必須 | |
| string | 必須 | メッセージをストリーミングするAWS Kinesisデータストリーム。 |
| string | 条件付き | 指定されたストリームが存在するAWSリージョン。Kinesis は、異なるリージョンで同じ名前の複数のデータ ストリームをサポートしています。同じ接続内の 2 つ以上のリージョンのデータストリームに同じ名前を使用する場合は、使用する名前とリージョンの組み合わせを指定するために、このフィールドを使用する必要があります。 |
| ドキュメント | 任意 | 受信メッセージの権限のあるタイムスタンプを定義するドキュメント。
|
| string | 任意 | プロジェクションされたドキュメントのタイムスタンプのフィールド名。このフィールドを使用して、デフォルトのタイムスタンプフィールド名を上書きします。 |
| ドキュメント | 任意 | 埋め込みの計算で無視される前に、シャードがアイドル状態になる時間を指定するドキュメント。 このフィールドは、デフォルトで無効になっています。アイドル状態であるため、前に移動しないシャードを取り扱うには、このフィールドに値を設定します。 |
| ドキュメント | 任意 | シャード アイドル タイムアウトの期間を指定する数値。 |
| ドキュメント | 任意 | シャード アイドル タイムアウトの期間の単位。
|
| ドキュメント | 任意 | のさまざまなデフォルト値を上書きするフィールドを含むドキュメント。 |
| string | 任意 | Kinesis データ ストリームの履歴内の、メッセージの取り込みを開始する位置。次のいずれかである必要があります。
デフォルトは |
| date | 条件付き | メッセージの取り込みを開始するタイムスタンプ。 |
| integer | 任意 | リシャーディングの目的でKinesisを超えるデータフローの速度をチェックする間隔(秒単位)。 デフォルトは |
ドキュメント配列
ドキュメントの配列を操作するために、 $sourceステージには次のプロトタイプ形式があります。
{ "$source": { "timeField": { $toDate | $dateFromString: <expression> }, "documents" : [{source-doc},...] | <expression> } }
$sourceステージは、次のフィールドを持つドキュメントを取得します。
フィールド | タイプ | 必要性 | 説明 |
|---|---|---|---|
| ドキュメント | 任意 | 受信メッセージの権限のあるタイムスタンプを定義するドキュメント。
|
| 配列 | 条件付き | ストリーミング データソースとして使用するドキュメントの配列。 このフィールドの値は、オブジェクトの配列、またはオブジェクトの配列として評価される 式 のいずれかになります。 |
CRON スケジュール
接続からの読み取りではなく、定期的なスケジュールでドキュメントを生成する場合、$source ステージのプロトタイプ形式は次のとおりです。
{ "$source": { "schedule": "<cron-expression>", "tsFieldName": "<timestamp-field-name>" } }
$sourceステージは、次のフィールドを持つドキュメントを取得します。
フィールド | タイプ | 必要性 | 説明 |
|---|---|---|---|
| string | 必須 | ステージがドキュメントをいつ発行するかを決定する 6 フィールドの CR式。フィールドは、秒、分、時間、日付、月、曜日の順に並べられます。 Atlas Stream Processing は UTC で式を評価します。 各フィールドは、すべての値、単一の値、 Atlas Stream Processing が、日付と曜日の両方を制限する式を解決する方法の詳細については、「 動作 」を参照してください。 |
| string | 任意 | ステージがスケジュールされたタイムスタンプをプロジェクションするフィールドの名前。デフォルトは |
動作
$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 の請求参照 」を参照してください。
例
Kafka の例
ストリーミング データソースは、気象用サンプル データセットのスキーマに準拠して、さまざまなロケーションから詳細な気象レポートを生成します。次の集計には 3 つのステージがあります。
$sourceステージでは、 という名前のトピックでこれらのレポートを収集するApachemy_weatherdataKafkaプロバイダーとの接続を確立し、各レコードが後続の集計ステージに取り込まれる際に公開します。このステージではまた、プロジェクションを実行するタイムスタンプフィールドの名前が上書きされ、ingestionTimeに設定されます。$matchステージでは、dewPoint.value5.0が 以下のドキュメントを除外し、 がdewPoint.value5.0より大きいドキュメントを次のステージに渡します。$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 コレクションに対する変更ストリームを開き、変更を記録します。
$sourceステージはcluster0-collectionソースに接続し、sample_weatherdataデータベース内のdataコレクションに対して変更ストリームを開きます。$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コレクションへの変更をレコード。
$sourceステージはdb-change-stream-connectionソースに接続し、sample_mflixソースデータベースに対して変更ストリームを開きます。config.startAtOperationTimeフィールドは、ソースがレポートを開始する時間を設定します。この例では、startAtの値を 1 分前に開始するように設定します。$mergeステージは、変更ストリームドキュメントをdb_changessample_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 が省略されているため、単一のコレクションやデータベースではなく、ソースクラスター上のすべてのデータベースとコレクションからの変更が報告されます。
$sourcecluster-changestream-connectionステージは ソースに接続し、ソースクラスター全体に対して変更ストリームを開きます。 フィールドは、ストリームconfig.startAtOperationTimeプロセッサが または 指定した時間の後に発生する変更のレポートを開始することを指定します。$mergeステージでは、宛先クラスターのeventscluster_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 つのステージを実行します。
$sourceステージでは、気象測定のインラインdocuments配列をストリーミングデータソースとして定義し、timeFieldを使用して各ドキュメントのtimestampフィールドを認証タイムスタンプとして指定します。$matchステージでは、dewPoint.valueが5.0を超えるドキュメントのみが次のステージに渡されます。$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データベース内の という名前のコレクションに書き込みます。
$sourcesample_stream_solarステージは ソースに接続します。 フィールドは、timeFieldtimestampを使用して各受信レポートの$dateFromStringフィールドを日付に変換します。$matchステージでは、device_idがdevice_8であるドキュメントを除外し、他のすべてのデバイスからのレポートを次のステージに渡します。ステージは、
$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 に書き込む前にそれらのレポートを除外するためです。
注意
前述の例はその一般的な例です。 ストリーミング データは静的ではなく、各ユーザーに異なるドキュメントが表示されます。
CRON スケジュールの例
次の集計では、接続からの読み取りではなく、5 分ごとにドキュメントが生成されます。この集計は3 つのステージを実行します。
$sourceステージでは05 分ごとの秒 に空のドキュメントが生成されます(UTC)。$projectステージでは、各ドキュメントにジョブ名がラベル付けされ、_tsフィールドから フィールドにスケジュールされたタイムスタンプがコピーされます。runAt$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()
注意
前の は の の例です。表示されるタイムスタンプは、ストリーム プロセッサを起動するときに異なります。