Definição
O estágio $tumblingWindow especifica uma janela em cascata para agregação de dados. As janelas do Atlas Stream Processing são com estado, podem ser recuperadas se interrompidas e têm mecanismos para processar dados atrasados. Você deve aplicar todas as outras queries de agregação aos seus dados de streaming dentro desse estágio da janela.
$tumblingWindowUm estágio de pipeline do
$tumblingWindowtem a seguinte forma de protótipo:{ "$tumblingWindow": { "boundary": "eventTime" | "processingTime", "interval": { "size": <int>, "unit": "<unit-of-time>" }, "pipeline" : [ <aggregation-stage-array> ], "offset": { "offsetFromUtc": <int>, "unit": "<unit-of-time>" }, "idleTimeout": { "size": <int>, "unit": "<unit-of-time>" }, "allowedLateness": { size: <int>, unit: "<unit-of-time>" } } } Como alternativa, uma etapa de pipeline
$tumblingWindowpode ter os camposallowedLatenesseidleTimeoutcom um valor inteiro de 0, conforme demonstrado abaixo:{ "$tumblingWindow": { "boundary": "eventTime" | "processingTime", "interval": { "size": <int>, "unit": "<unit-of-time>" }, "pipeline" : [ <aggregation-stage-array> ], "offset": { "offsetFromUtc": <int>, "unit": "<unit-of-time>" }, "idleTimeout": 0, "allowedLateness": 0 } }
Sintaxe
O estágio $tumblingWindow recebe um documento com os seguintes campos:
Campo | Tipo | necessidade | Descrição |
|---|---|---|---|
| string | Opcional | String que especifica se os limites da janela são determinados pelo tempo de evento ou pelo tempo de processamento. O valor pode ser
|
| documento | Obrigatório | Documento que especifica o intervalo de uma janela de tombamento como uma combinação de um tamanho e uma unidade de tempo em que:
Por exemplo, um |
| array | Obrigatório | Pipeline de agregação aninhado avaliado em relação às mensagens dentro da janela. |
| documento | Opcional | Documento que especifica uma compensação de horário para os limites da janela em relação ao UTC. O documento é uma combinação do campo de tamanho
Por exemplo, um |
| documento | Opcional | Documento que especifica quanto tempo esperar antes de fechar Windows se
Se você definir Por exemplo, considere uma janela 12:00 pm a 1:00 pm e Alternativamente, você pode definir esta configuração com um valor inteiro de.0 Consulte a definição do pipeline para obter mais informações. Para ver como ocorre a interação de um tempo limite de inatividade com as marcas d'água ao longo de uma janela, consulte Tempo limite de inatividade. |
| documento | Opcional | Documento que especifica por quanto tempo manter abertas as janelas geradas a partir da fonte para aceitar dados que chegam tarde após o processamento dos documentos até o horário de término da janela. Se omitido, o padrão é 3 segundos. Alternativamente, você pode definir esta configuração com um valor inteiro de.0 Consulte a definição do pipeline para obter mais informações. Para ver como o atraso permitido atrasa o fechamento de uma janela em relação à marca d'água, consulte Atraso permitido. |
Comportamento
O Atlas Stream Processing permite apenas um estágio de janela por pipeline.
Quando você aplica o $group estágio ao estágio da janela, uma única chave de grupo tem um limite de 100 megabytes de RAM.
A compatibilidade para determinados estágios de agregação pode ser limitada ou indisponível nas janelas. Para saber mais, consulte Estágios de pipeline de agregação com compatibilidade.
No caso de uma interrupção de serviço, você pode retomar o pipeline interno de uma janela a partir do estado em que se encontrava no momento da interrupção. Para saber mais, consulte checkpoints.
Exemplos
Uma fonte de dados de streaming gera relatórios meteorológicos detalhados a partir de vários locais, em conformidade com o esquema do Conjunto de dados meteorológicos de amostra. A seguinte agregação tem três estágios:
O estágio estabelece uma
$sourceconexão com o intermediário do Apache Kafka, coletando esses relatórios em um tópicomy_weatherdatachamado, expondo cada registro à medida que ele é ingerido para os estágios de agregação subsequentes.O estágio define janelas consecutivas
$tumblingWindowcom 30duração de segundos. Cada janela executa umpipelineinterno, que encontra a média, a mediana, o máximo e o mínimoatmosphericPressureObservation.altimeterSetting.valuede para a duração dessa janela. Empipelineseguida, o gera um único documento com um_idequivalente ao carimbo de data/hora de início da janela que representa e os valores especificados para essa janela.O estágio escreve o resultado em uma collection do Atlas
$mergechamadastreamnosample_weatherstreambanco de dados. Se não existir tal banco de dados ou coleção, o Atlas os criará.
{ '$source': { connectionName: 'sample_weatherdata', topic: 'my_weatherdata', tsFieldName: 'ingestionTime' } }, { '$tumblingWindow': { interval: { size: 30, unit: "second" }, pipeline: [{ $group: { _id: { $meta: "stream.window.start"}, averagePressure: { $avg: "$atmosphericPressureObservation.altimeterSetting.value" }, medianPressure: { $median: { input: "$atmosphericPressureObservation.altimeterSetting.value", method: "approximate" } }, maxPressure: { $max: "$atmosphericPressureObservation.altimeterSetting.value" }, minPressure: { $min: "$atmosphericPressureObservation.altimeterSetting.value" } } }] } }, { '$merge': { into: { connectionName: 'weatherStreamOutput', db: 'sample_weatherstream', coll: 'stream' } } }
Para visualizar os documentos na coleção sample_weatherstream.stream resultante, conecte-se ao cluster Atlas e execute o seguinte comando:
db.getSiblingDB("sample_weatherstream").stream.find()
{ _id: ISODate('2024-09-26T16:34:00.000Z'), averagePressure: 5271.47894736842, maxPressure: 9999.9, medianPressure: 1015.9, minPressure: 1015.9 }, { _id: ISODate('2024-09-26T16:34:30.000Z'), averagePressure: 5507.9, maxPressure: 9999.9, medianPressure: 1015.9, minPressure: 1015.9 }
Observação
O exemplo anterior é representativo. Os dados de streaming não são estáticos e cada usuário vê documentos distintos.