You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark/Scala:处理每个设备最近10条非空观测数据需求

Got it, let's tackle this problem head-on. You need to generate a daily sensor data report that focuses on the last 10 measurements per item, while handling edge cases like items with no daily data, multi-day gaps, and newly added sensors. Since you're using a Scala case class, I'll frame this with idiomatic Scala (and touch on Spark if you're dealing with large-scale data).

1. Core Logic: Get Last 10 Measurements per Item

First, let's start with the basic grouping and filtering. The key here is to sort measurements by date (descending) and take the top 10 for each item.

Pure Scala Implementation

If you're working with in-memory data:

case class Measurement(date: String, item: String, value: Int)

// Sample data (mix of new items, existing items with gaps)
val rawData = List(
  Measurement("2024-05-20", "sensor_01", 23),
  Measurement("2024-05-19", "sensor_01", 22),
  Measurement("2024-05-10", "sensor_01", 25), // Gap between 10 and 19
  Measurement("2024-05-20", "sensor_02", 45),
  Measurement("2024-05-20", "sensor_03", 10), // New item today
  // sensor_04 has no data today
)

// Step 1: Group measurements by item
val groupedByItem = rawData.groupBy(_.item)

// Step 2: For each item, sort by date descending and take last 10
val last10PerItem = groupedByItem.map { case (item, measurements) =>
  val sorted = measurements.sortBy(_.date)(Ordering[String].reverse)
  (item, sorted.take(10))
}

Spark Implementation (for large datasets)

If you're dealing with big data, Spark's window functions are perfect for this:

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

// Assume you have a DataFrame of Measurement data
val windowSpec = Window.partitionBy("item").orderBy(desc("date"))

val last10DF = rawDataDF
  .withColumn("row_num", row_number().over(windowSpec))
  .filter(col("row_num") <= 10)
  .drop("row_num")
2. Handling Missing & New Items

The tricky part is ensuring items with no daily data (or multi-day gaps) are still included in the report, using their most recent 10 measurements. Also, new items should appear as soon as they have data.

Track All Known Items

To handle missing items, you need a way to keep track of all items that have ever been recorded. You can store this list in a persistent store (like a database, file, or even a Spark checkpoint) between daily runs.

For example:

// Let's say we load the full list of known items from a previous run
val allKnownItems = List("sensor_01", "sensor_02", "sensor_03", "sensor_04")

// For items not in today's data, we need to retain their last 10 measurements from history
// (You'd store this historical data persistently, e.g., in a parquet file)
val historicalLast10 = Map(
  "sensor_04" -> List(
    Measurement("2024-05-18", "sensor_04", 30),
    Measurement("2024-05-17", "sensor_04", 31)
  )
)

// Combine today's last 10 with historical data for missing items
val fullReport = allKnownItems.map { item =>
  last10PerItem.getOrElse(item, historicalLast10.getOrElse(item, List.empty))
}.filter(_.nonEmpty) // Optional: exclude items with zero historical data

Handling New Items

New items will automatically be added to allKnownItems once they appear in the daily data. After processing each day, update your persistent list of known items:

val updatedKnownItems = (allKnownItems ++ rawData.map(_.item)).distinct
// Save updatedKnownItems to your persistent store for tomorrow's run
3. Generate the Daily Report

Finally, format the data into a readable report. For example:

def generateReport(reportData: List[(String, List[Measurement])]): String = {
  reportData.map { case (item, measurements) =>
    s"## Item: $item\n" +
    measurements.map(m => s"- ${m.date}: ${m.value}").mkString("\n") +
    "\n"
  }.mkString("\n")
}

println(generateReport(last10PerItem.toList))

This will output a clean markdown-style report that includes every known item, even those with no new data that day.


内容的提问来源于stack exchange,提问作者Averell

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.25 06:13:27