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).
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")
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
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

