如何在Spark Scala中实现DataFrame新行与历史行的合并?
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
- Convert Dates: We first turn the string dates into proper
DateTypecolumns so we can easily compare which dates are earlier. - Get Unique Dates: We extract all distinct dates from the input data and sort them chronologically.
- 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 beforecurrent_date. This ensures eachcurrent_dateincludes all historical data up to that point. - Format & Select: We convert
current_dateback 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

