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

如何基于Apache Spark和Sparkube实现‘Last Non Empty’度量?

Nice question! Handling "Last Non Empty" aggregations for time-series data like inventory is a common use case, and since Sparkube doesn’t support this natively, we can precompute the required values using Apache Spark before feeding the data into Sparkube. Let’s walk through the solution step by step.

Step 1: Clarify the Requirement

For each product and time granularity (like year), we need the stock value from the most recent record in that period—not a sum, average, or other standard aggregation. For example, 2018's inventory for Oranges should be 46000 (from the 2018-03-01 record), not the total of all 2018 entries.

Step 2: Preprocess Data in Spark

We'll use Spark to filter down to only the latest valid record per product and time group. Here are two reliable approaches:

Approach 1: Group by Dimensions + Join Back

This method first identifies the latest date for each product-year pair, then joins back to the original dataset to fetch the corresponding stock value:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, year, max as spark_max

# Initialize Spark session
spark = SparkSession.builder.appName("LastNonEmptyInventory").getOrCreate()

# Load your dataset (replace with your actual data source like S3, HDFS, etc.)
inventory_data = [
    ("2017-11-01", "Oranges", 40000),
    ("2017-11-01", "Apples", 120000),
    ("2017-12-01", "Oranges", 42000),
    ("2017-12-01", "Apples", 110000),
    ("2018-01-01", "Oranges", 50000),
    ("2018-01-01", "Apples", 100000),
    ("2018-02-01", "Oranges", 48000),
    ("2018-02-01", "Apples", 130000),
    ("2018-03-01", "Oranges", 46000),
    ("2018-03-01", "Apples", 120000)
]

df = spark.createDataFrame(inventory_data, ["Time", "Product", "Stock"])
# Convert Time to date type and extract Year for grouping
df = df.withColumn("Time", col("Time").cast("date"))
df = df.withColumn("Year", year(col("Time")))

# Get the latest date for each Product + Year combination
latest_dates = df.groupBy("Product", "Year").agg(spark_max("Time").alias("LatestTime"))

# Join back to the original data to get the stock value for that latest date
last_non_empty_df = latest_dates.join(
    df,
    on=["Product", "Time"],
    how="inner"
).select("Product", "Year", "Time", "Stock")

# View the final preprocessed data
last_non_empty_df.show()

Running this code will output exactly the records we need:

+-------+----+----------+-----+
|Product|Year|      Time|Stock|
+-------+----+----------+-----+
|Oranges|2017|2017-12-01|42000|
| Apples|2017|2017-12-01|110000|
|Oranges|2018|2018-03-01|46000|
| Apples|2018|2018-03-01|120000|
+-------+----+----------+-----+

Approach 2: Window Functions (For Multiple Granularities)

If you need to support both yearly and monthly "Last Non Empty" values, window functions are more flexible. We'll rank records by time descending within each dimension group, then keep only the top-ranked (latest) entry:

from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, month

# Add Month column for finer-grained grouping
df = df.withColumn("Month", month(col("Time")))

# Window for yearly latest records: partition by Product + Year, order by Time descending
year_window = Window.partitionBy("Product", "Year").orderBy(col("Time").desc())
df_year_ranked = df.withColumn("year_rank", row_number().over(year_window))
last_non_empty_year = df_year_ranked.filter(col("year_rank") == 1).drop("year_rank")

# Window for monthly latest records
month_window = Window.partitionBy("Product", "Year", "Month").orderBy(col("Time").desc())
df_month_ranked = df.withColumn("month_rank", row_number().over(month_window))
last_non_empty_month = df_month_ranked.filter(col("month_rank") == 1).drop("month_rank")

Step 3: Use Preprocessed Data in Sparkube

Import your preprocessed dataset (like last_non_empty_df) into Sparkube. Since each Product + Year (or Product + Year + Month) pair has exactly one record, you can use any of Sparkube's built-in aggregations (SUM, MAX, AVG, etc.) to get the "Last Non Empty" value—they'll all return the single valid stock value for that group.

Why This Works

Sparkube doesn’t natively support custom aggregations like "Last Non Empty", but by precomputing the latest valid records in Spark, we simplify the problem to aggregating over unique dimension pairs. This plays perfectly to Sparkube's strengths in building fast, interactive multidimensional cubes.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:22:07