정의
현재 스트림 처리 작업 공간에 임시 스트림 프로세서를 생성합니다.
호환성
이 메서드는 Atlas Stream Processing Workspaces에서 지원됩니다.
구문
sp.process() 메서드의 구문은 다음과 같습니다.
sp.process( [ <pipeline> ] )
명령 필드
sp.process() 는 다음 필드를 사용합니다.
필드 | 유형 | 필요성 | 설명 |
|---|---|---|---|
| 배열 | 필수 사항 | 스트림 집계 파이프라인 스트림 프로세서는 스트리밍 데이터에 적용됩니다. |
| 객체 | 옵션 | 스트림 프로세서의 선택적 설정을 정의하는 객체입니다. |
| integer | 옵션 | 터미널로 반환할 수 있는 최대 문서 수입니다. |
| 문자열 | 옵션 | 스트림 프로세서에 할당할 계층입니다. 계층 이름은 대소문자를 구분하지 않습니다. 이 필드 생략하면 스트림 프로세서는 스트림 처리 작업 공간에 기본값 계층 사용합니다. 다음 값 중 하나여야 합니다.
|
행동
sp.process() 현재 스트림 처리 작업 공간에 이름이 없는 임시 스트림 프로세서를 생성하고 즉시 초기화합니다. 이 스트림 프로세서는 실행되는 동안에만 유지됩니다. 임시 스트림 프로세서를 종료한 경우 다시 생성해야 사용 가능합니다.
임시 스트림 프로세서는 다음 중 하나라도 발생하면 중지됩니다.
프로세서는
options.limit에서 설정하다 문서 수를 반환합니다.10 minutes elapse.
프로세서를 중지합니다.
액세스 제어
를 실행 sp.process() 사용자는 atlasAdmin 역할 있어야 합니다.
예시
다음 예시 연결에서 데이터를 수집하는 임시 스트림 프로세서를 sample_stream_solar 만듭니다. 프로세서는 필드 값이 인 모든 문서를 제외하고 device_id device_8 나머지는 10지속 시간이 초인 텀블링 창 으로 전달합니다. 각 창 수신한 문서를 그룹화한 다음 각 그룹 의 다양한 통계를 반환합니다. 그런 다음 스트림 프로세서는 solar_db.solar_coll mongodb1 연결을 통해 이러한 레코드를 에 병합합니다.
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"] } } ] )