0Pricing
MongoDB Academy · 课时

对时间序列执行窗口聚合

您将使用带有基于时间的窗口边界的 $setWindowFields,计算传感器读数的滚动平均值和累计总和。

对时间序列执行窗口聚合 是 CoddyKit 上的免费 MongoDB Academy 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 MongoDB Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 MongoDB Academy 课程共包含 4 节课。

本课时的部分内容尚未翻译,以英文显示。

What Are Window Functions?

A window function computes a value for each document in a result set by looking at a surrounding range (the 'window') of documents — without collapsing them into a group the way $group does. Each input document produces exactly one output document, but the computed value reflects aggregation across the window. MongoDB added window functions via the $setWindowFields stage in version 5.0.

The $setWindowFields Stage

$setWindowFields is the aggregation stage that enables window functions in MongoDB. It accepts a partitionBy expression (similar to SQL's PARTITION BY), an output object defining new fields and their window operators, and a sortBy document that establishes the ordering within each partition. This combination makes it ideal for time series: partition by sensor, sort by timestamp, compute running totals or moving averages.

db.sensorReadings.aggregate([
  {
    $setWindowFields: {
      partitionBy: '$sensorId',    // one window per sensor
      sortBy: { timestamp: 1 },   // ordered by time ascending
      output: {
        runningTotal: {
          $sum: '$temperature',
          window: { documents: ['unbounded', 'current'] }
        }
      }
    }
  }
])

Document-Based Window Boundaries

Window boundaries can be specified in terms of document positions relative to the current document. The special keywords 'unbounded' (from the first document in the partition) and 'current' (the current document) define common boundaries. You can also use integers: [-3, 0] means the 3 preceding documents and the current one — perfect for a rolling 4-point moving average.

// 5-point rolling average (2 before, current, 2 after)
db.sensorReadings.aggregate([
  {
    $setWindowFields: {
      partitionBy: '$sensorId',
      sortBy: { timestamp: 1 },
      output: {
        rollingAvgTemp: {
          $avg: '$temperature',
          window: { documents: [-2, 2] }
        }
      }
    }
  }
])

Range-Based Time Windows

For time series data, range-based windows are more natural than document-based ones because data may not arrive at uniform intervals. Range windows use a unit (milliseconds by default) and an amount relative to the current document's sort value. For example, a 1-hour trailing window would be { range: [-3600000, 0], unit: 'millisecond' }. Dedicated time units like 'hour', 'minute', and 'second' are supported too.

// 1-hour trailing moving average by time range
db.sensorReadings.aggregate([
  {
    $setWindowFields: {
      partitionBy: '$sensorId',
      sortBy: { timestamp: 1 },
      output: {
        movingAvgTemp: {
          $avg: '$temperature',
          window: {
            range: [-1, 0],
            unit: 'hour'
          }
        }
      }
    }
  }
])

Running Totals With $sum

A running total (or cumulative sum) is computed by setting the window from 'unbounded' to 'current'. Each document's output field reflects the sum of all values from the first document in the partition up to and including the current one. This is useful for computing cumulative energy consumption, total transactions, or bytes transferred over time.

// Cumulative energy consumption per device
db.energyReadings.aggregate([
  {
    $setWindowFields: {
      partitionBy: '$deviceId',
      sortBy: { timestamp: 1 },
      output: {
        cumulativeKwh: {
          $sum: '$kwh',
          window: { documents: ['unbounded', 'current'] }
        }
      }
    }
  },
  { $project: { _id: 0, timestamp: 1, deviceId: 1, kwh: 1, cumulativeKwh: 1 } }
])

Ranking With $rank and $denseRank

$rank and $denseRank are window operators that assign position numbers to documents within their partition, ordered by the sortBy expression. $rank leaves gaps in numbering for ties (1, 2, 2, 4…) while $denseRank does not (1, 2, 2, 3…). These are useful for ranking sensors by their latest reading or identifying top-N performing devices.

// Rank sensors by their average temperature (after grouping)
db.sensorStats.aggregate([
  {
    $setWindowFields: {
      partitionBy: null,         // no partition — global rank
      sortBy: { avgTemp: -1 },  // highest temp ranked first
      output: {
        rank: { $rank: {} }
      }
    }
  }
])

Row Numbers With $documentNumber

$documentNumber assigns a sequential integer starting at 1 to each document within its partition, ordered by sortBy. Unlike $rank, it never produces ties — every document gets a unique number. This is handy for pagination, batch numbering of exports, or labelling sequential events in a sensor stream.

db.sensorReadings.aggregate([
  {
    $setWindowFields: {
      partitionBy: '$sensorId',
      sortBy: { timestamp: 1 },
      output: {
        readingNumber: { $documentNumber: {} }
      }
    }
  },
  { $match: { sensorId: 'sensor-42' } }
])

Shifting Values: $shift Operator

The $shift operator accesses the value of a field from a document at a relative position. { $shift: { output: '$temperature', by: -1 } } returns the temperature from the previous document in the partition. This is perfect for computing deltas — the difference between the current reading and the one before it — which reveals rate-of-change in sensor data.

// Compute temperature delta vs previous reading
db.sensorReadings.aggregate([
  {
    $setWindowFields: {
      partitionBy: '$sensorId',
      sortBy: { timestamp: 1 },
      output: {
        prevTemp: {
          $shift: { output: '$temperature', by: -1, default: null }
        }
      }
    }
  },
  {
    $addFields: {
      tempDelta: { $subtract: ['$temperature', '$prevTemp'] }
    }
  }
])

Combining $setWindowFields With $match

Always place a $match stage before $setWindowFields to reduce the working set. Window computations happen in memory, so feeding fewer documents into the stage significantly reduces RAM usage. After the window stage, you can add another $match to filter the enriched output — for example, keeping only documents where the rolling average exceeds a threshold.

// Flag anomalies where temp spikes above rolling average
db.sensorReadings.aggregate([
  {
    $match: {
      sensorId: 'sensor-42',
      timestamp: { $gte: new Date('2024-06-01T00:00:00Z') }
    }
  },
  {
    $setWindowFields: {
      partitionBy: '$sensorId',
      sortBy: { timestamp: 1 },
      output: {
        rollingAvg: {
          $avg: '$temperature',
          window: { documents: [-9, 0] }
        }
      }
    }
  },
  {
    $match: {
      $expr: { $gt: ['$temperature', { $multiply: ['$rollingAvg', 1.2] }] }
    }
  }
])

Performance: Partition Size Matters

$setWindowFields must hold each partition in memory to compute window values. For time series data, this means one sensor's data for the requested time range. If your partitions are extremely large (millions of measurements), consider pre-filtering with $match or using allowDiskUse: true on the aggregation to spill partitions to disk. Alternatively, pre-aggregate into per-hour summaries before applying window functions.

// Allow disk use for very large partitions
db.sensorReadings.aggregate(
  [
    { $match: { timestamp: { $gte: new Date('2024-01-01T00:00:00Z') } } },
    {
      $setWindowFields: {
        partitionBy: '$sensorId',
        sortBy: { timestamp: 1 },
        output: {
          rolling7d: {
            $avg: '$temperature',
            window: { range: [-7, 0], unit: 'day' }
          }
        }
      }
    }
  ],
  { allowDiskUse: true }
)

Practical Example: Anomaly Detection

A complete anomaly detection pipeline combines $match for range filtering, $setWindowFields with a trailing window to compute the rolling average and standard deviation, and a final $match with $expr to surface readings that deviate more than 2 standard deviations from the local mean. This pattern forms the backbone of real-time IoT monitoring systems.

db.sensorReadings.aggregate([
  { $match: { timestamp: { $gte: new Date('2024-06-01T00:00:00Z') } } },
  {
    $setWindowFields: {
      partitionBy: '$sensorId',
      sortBy: { timestamp: 1 },
      output: {
        rollingAvg:    { $avg: '$temperature', window: { documents: [-19, 0] } },
        rollingStdDev: { $stdDevSamp: '$temperature', window: { documents: [-19, 0] } }
      }
    }
  },
  {
    $match: {
      $expr: {
        $gt: [
          { $abs: { $subtract: ['$temperature', '$rollingAvg'] } },
          { $multiply: ['$rollingStdDev', 2] }
        ]
      }
    }
  }
])

Quick Check

Test your understanding of MongoDB & NoSQL Databases concepts from this lesson.

Lesson Recap

In this lesson you learned: $setWindowFields enriches each document with values computed over a surrounding window without collapsing the result, document and range boundaries give you fine-grained control over trailing, leading, or time-span windows, and operators like $avg, $sum, $shift, and $rank cover moving averages, running totals, delta computation, and ranking. Next up we explore automatic data expiration using expireAfterSeconds.

常见问题解答

「对时间序列执行窗口聚合」课时是免费的吗?

是的 — 「对时间序列执行窗口聚合」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 MongoDB Academy 课程的其余内容,请升级到 CoddyKit PRO。 MongoDB Academy 课程共包含 4 节课。

「对时间序列执行窗口聚合」这节课中我会学到什么?

您将使用带有基于时间的窗口边界的 $setWindowFields,计算传感器读数的滚动平均值和累计总和。 你通过在浏览器中直接运行的动手代码来练习 MongoDB Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 MongoDB Academy 需要有经验吗?

无需任何先前经验。CoddyKit 上的 MongoDB Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 3 节课,共 4 节。

「对时间序列执行窗口聚合」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 MongoDB Academy 课中编写并运行代码吗?

能。每节 MongoDB Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 创建时间序列集合
  2. 插入和查询时间序列数据
  3. 对时间序列执行窗口聚合
  4. 使用 expireAfterSeconds 自动过期数据
← 返回 MongoDB Academy