정의
sp.process()현재 Stream Processing Workspace 에 임시 Stream Processor 를 생성합니다.
stream processing 작업 공간에 연결된 동안에만 이 명령을 호출할 수 있습니다.
이 명령에는
mongosh버전 이상 2.0 이(가) 필요합니다.
호환성
이 메서드는 Atlas Stream Processing Workspaces에서 지원됩니다.
구문
sp.process() 메서드의 구문은 다음과 같습니다.
sp.process( [ <pipeline> ], { <options> } )
명령 필드
sp.process() 는 다음 필드를 사용합니다.
필드 | 유형 | 필요성 | 설명 |
|---|---|---|---|
| 배열 | 필수 사항 | 스트리밍 데이터에 적용되는 스트리밍 프로세서의 스트리밍 집계 파이프라인. |
| 객체 | 옵션 | 스트림 프로세서에 대한 선택 사항을 정의하는 객체입니다. |
| integer | 옵션 | 터미널로 반환할 문서의 최대 개수입니다. |
| 문자열 | 옵션 | Stream Processing에 할당할 계층입니다. 계층 이름은 대소문자를 구분하지 않습니다. 이 필드를 생략하면 Stream Processing Workspace의 기본 계층을 사용합니다. 다음 값 중 하나여야 합니다.
|
행동
sp.process() 현재 Stream Processing 작업 공간에 임시적인 이름 없는 스트림 프로세서를 생성하고 즉시 초기화합니다. 이 Stream Processor는 실행 중에만 유지됩니다. 임시 Stream Processor를 종료하는 경우 사용하려면 다시 생성해야 합니다.
일시적 스트림 프로세서는 다음 중 하나라도 발생하면 중지됩니다.
프로세서는
options.limit에 설정된 문서 수를 반환합니다.10 minutes elapse.
프로세서를 중지합니다.
액세스 제어
를 실행 sp.process() 사용자는 atlasAdmin 역할 있어야 합니다.
예시
다음 예시는 sample_stream_solar 연결에서 데이터를 수집하는 일시적 스트림 프로세서를 만듭합니다. 프로세서는 device_id 필드의 값이 device_8인 모든 문서를 제외하고 나머지는 10초 기간의 터블링 창으로 전달합니다. 각 창은 수신한 문서를 그룹화한 다음 각 그룹의 다양한 통계를 반환합니다. 스트림 프로세서는 이러한 기록을 병합하여 mongodb1 연결을 통해 solar_db.solar_coll 합니다.
sp.process( [ { $source: { connectionName: 'sample_stream_solar', timeField: { $dateFromString: { dateString: '$timestamp' } } } }, { $match: { $expr: { $ne: [ "$device_id", "device_8" ] } } }, { $tumblingWindow: { interval: { size: Int32(10), unit: "second" }, "pipeline": [ { $group: { "_id": { "device_id": "$device_id" }, "max_temp": { $max: "$obs.temp" }, "max_watts": { $max: "$obs.watts" }, "min_watts": { $min: "$obs.watts" }, "avg_watts": { $avg: "$obs.watts" }, "median_watts": { $median: { input: "$obs.watts", method: "approximate" } } } } ] } }, { $merge: { into: { connectionName: "mongodb1", db: "solar_db", coll: "solar_coll" }, on: ["_id"] } } ] )