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

如何在SQL及Scala-Spark中基于条件引用跨行指定列值

跨年份销售数据填充解决方案

SQL 查询语句

假设数据表为product_sales,包含字段product_id、year、this_year_sales、last_year_sales,以下两种方式均可实现需求:

方式1:自连接实现

适合简单的跨表关联场景,可直接查询或更新原表:

-- 查询2023年数据并填充去年销售额
SELECT
    curr.product_id,
    curr.year,
    curr.this_year_sales,
    COALESCE(prev.this_year_sales, 0) AS last_year_sales
FROM product_sales curr
LEFT JOIN product_sales prev
    ON curr.product_id = prev.product_id
    AND curr.year = prev.year + 1
WHERE curr.year = 2023;

-- 更新原表2023年的last_year_sales字段
UPDATE product_sales curr
SET last_year_sales = (
    SELECT prev.this_year_sales
    FROM product_sales prev
    WHERE prev.product_id = curr.product_id
    AND prev.year = curr.year - 1
)
WHERE curr.year = 2023;

方式2:窗口函数(LAG)实现

更高效的同组内前后行关联,适合处理大规模数据:

-- 查询2023年数据并生成去年销售额
SELECT
    product_id,
    year,
    this_year_sales,
    LAG(this_year_sales, 1) OVER (
        PARTITION BY product_id
        ORDER BY year
    ) AS last_year_sales
FROM product_sales
WHERE year = 2023;

-- 以MySQL为例,更新原表字段
WITH sales_with_lag AS (
    SELECT
        product_id,
        year,
        LAG(this_year_sales, 1) OVER (
            PARTITION BY product_id
            ORDER BY year
        ) AS lag_sales
    FROM product_sales
)
UPDATE product_sales curr
JOIN sales_with_lag sl
    ON curr.product_id = sl.product_id
    AND curr.year = sl.year
SET curr.last_year_sales = sl.lag_sales
WHERE curr.year = 2023;

Scala-Spark 实现方案

假设已有Spark DataFrame originalDF,结构为product_id(String/Int)、year(Int)、this_year_sales(Double):

1. 导入依赖包

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

2. 计算并填充去年销售额

// 定义窗口:按产品分组,按年份升序排序
val productWindow = Window.partitionBy("product_id").orderBy("year")

// 生成last_year_sales字段,仅保留2023年数据
val resultDF = originalDF
    .withColumn("last_year_sales", lag("this_year_sales", 1).over(productWindow))
    .filter(col("year") === 2023)

// 如果需要保留所有年份数据并更新last_year_sales
val fullUpdatedDF = originalDF
    .withColumn("last_year_sales", lag("this_year_sales", 1).over(productWindow))

3. 输出或保存结果

// 打印结果预览
resultDF.show()

// 保存到目标表(示例为Hive表)
resultDF.write.mode("overwrite").saveAsTable("target_product_sales")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 03:42:56