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

如何在Spark Scala中实现DataFrame新行与历史行的合并?

Solution to Generate OLAP-style DataFrame from OLTP Data

Got it, let's break down how to solve this problem. The core requirement is to expand your OLTP data so that each date entry includes all historical rows from previous dates plus the new rows added on that date. Here's how to do it with Spark, using both Scala and Python examples, plus an optimized approach for large datasets.

Step 1: Prepare the Input Data

First, let's create our sample DataFrame and convert the string dates to actual date types (this makes date comparisons much easier).

Scala Code

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

// Sample OLTP data
val data = Seq(
  ("abc", "4/6/2018", 100),
  ("abc", "4/6/2018", 200),
  ("abc", "4/13/2018", 300)
)

val df = spark.createDataFrame(data).toDF("name", "date", "amount")
// Convert string date to date type for comparison
val dfWithDate = df.withColumn("date_dt", to_date(col("date"), "M/d/yyyy"))

Python Code

from pyspark.sql import functions as F

# Sample OLTP data
data = [
    ("abc", "4/6/2018", 100),
    ("abc", "4/6/2018", 200),
    ("abc", "4/13/2018", 300)
]

df = spark.createDataFrame(data, ["name", "date", "amount"])
# Convert string date to date type for comparison
df_with_date = df.withColumn("date_dt", F.to_date(F.col("date"), "M/d/yyyy"))

Step 2: Generate the OLAP DataFrame (Optimized for Large Datasets)

For big data, we want to avoid collecting data to the driver node. Instead, use a cross join with filtered dates:

Scala Optimized Code

// Get all unique dates, renamed as current_date
val allDates = dfWithDate.select("date_dt").distinct().orderBy("date_dt").withColumnRenamed("date_dt", "current_date")

// Cross join with original data, filter historical rows, and format the result date
val finalOlapDf = allDates.crossJoin(dfWithDate)
  .filter(col("date_dt") <= col("current_date"))
  .withColumn("date", date_format(col("current_date"), "M/d/yyyy"))
  .select("name", "date", "amount")
  .orderBy("date", "amount")

// Show the result
finalOlapDf.show()

Python Optimized Code

# Get all unique dates, renamed as current_date
all_dates = df_with_date.select("date_dt").distinct().orderBy("date_dt").withColumnRenamed("date_dt", "current_date")

# Cross join with original data, filter historical rows, and format the result date
final_olap_df = all_dates.crossJoin(df_with_date) \
  .filter(F.col("date_dt") <= F.col("current_date")) \
  .withColumn("date", F.date_format(F.col("current_date"), "M/d/yyyy")) \
  .select("name", "date", "amount") \
  .orderBy("date", "amount")

# Show the result
final_olap_df.show()

How This Works

  1. Convert Dates: We first turn the string dates into proper DateType columns so we can easily compare which dates are earlier.
  2. Get Unique Dates: We extract all distinct dates from the input data and sort them chronologically.
  3. Cross Join & Filter: We join each unique date (as current_date) with every row in the original data, then keep only rows where the original date is on or before current_date. This ensures each current_date includes all historical data up to that point.
  4. Format & Select: We convert current_date back to the original string format, then select the columns we need to match the expected OLAP structure.

Alternative Approach (Small Datasets)

If your dataset is small enough, you can collect the dates to the driver and iterate through them:

Scala Small Dataset Code

// Collect sorted unique dates to driver
val dateList = dfWithDate.select("date_dt").distinct().orderBy("date_dt").collect().map(_.getAs[java.sql.Date](0))

// Generate a DataFrame for each date, then union them all
val resultDfs = dateList.map { currentDate =>
  dfWithDate.filter(col("date_dt") <= currentDate)
    .withColumn("date", date_format(currentDate, "M/d/yyyy"))
    .select("name", "date", "amount")
}

val finalOlapDf = resultDfs.reduce(_ union _)
finalOlapDf.orderBy("date", "amount").show()

Python Small Dataset Code

# Collect sorted unique dates to driver
date_list = [row.date_dt for row in df_with_date.select("date_dt").distinct().orderBy("date_dt").collect()]

# Generate a DataFrame for each date, then union them all
result_dfs = []
for current_date in date_list:
    filtered_df = df_with_date.filter(F.col("date_dt") <= current_date)
    result_df = filtered_df.withColumn("date", F.date_format(current_date, "M/d/yyyy")) \
                          .select("name", "date", "amount")
    result_dfs.append(result_df)

final_olap_df = result_dfs[0]
for df in result_dfs[1:]:
    final_olap_df = final_olap_df.union(df)

final_olap_df.orderBy("date", "amount").show()

This will output exactly the OLAP structure you're expecting:

+----+---------+------+
|name|date     |amount|
+----+---------+------+
|abc |4/6/2018 | 100  |
|abc |4/6/2018 | 200  |
|abc |4/13/2018| 100  |
|abc |4/13/2018| 200  |
|abc |4/13/2018| 300  |
+----+---------+------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 10:07:42