SparkR SQL:基于时间戳在滚动时间窗口内统计符合条件的记录数
Got it, let's break down how to calculate rolling time window counts with SparkR SQL using your sample dataset. First, let's finish up that sample dataframe code since it was cut off—we'll add user IDs to make the grouping meaningful (since you probably want per-user rolling stats):
Complete Sample Dataset Code
# Generate timestamp sequence dates <- base::seq.POSIXt( from = as.POSIXlt(as.Date("2018-01-01"), format="%Y-%m-%d"), to = as.POSIXlt(as.Date("2018-01-03"), format="%Y-%m-%d"), by = "hour" ) # Status options and sampled values possible_statuses <- c('moving', 'stopped') statuses4demo <- base::sample(possible_statuses, size = length(dates)*2, replace = TRUE, prob = c(.75, .25)) # Final dataframe with 2 demo users df <- data.frame( date = rep(dates, 2), user_id = rep(c("user_1", "user_2"), each = length(dates)), status = statuses4demo )
Step 1: Set Up SparkR and Load Data
First, initialize your Spark session and convert the local R dataframe to a Spark DataFrame:
library(SparkR) # Start Spark session sparkR.session(appName = "RollingWindowStatusCount") # Convert local df to Spark DataFrame spark_df <- createDataFrame(df)
Step 2: Register a Temporary View for SQL Queries
To use Spark SQL, you need to register the Spark DataFrame as a temporary view (think of it as a temporary table you can query):
createOrReplaceTempView(spark_df, "device_status")
Step 3: Write the Rolling Window SQL Query
Let's say you want to calculate, for each user and timestamp, the number of 'moving' records in the previous 5 hours (a rolling window). Here's how to write that query:
SELECT user_id, date, status, -- Count only 'moving' records in the 5-hour rolling window COUNT(CASE WHEN status = 'moving' THEN 1 END) OVER ( PARTITION BY user_id -- Calculate window per user ORDER BY UNIX_TIMESTAMP(date) * 1000 -- Convert timestamp to milliseconds for range logic RANGE BETWEEN 5*3600*1000 PRECEDING AND CURRENT ROW -- Window = past 5h to current time ) AS rolling_5h_moving_count FROM device_status ORDER BY user_id, date
Run this query in SparkR and view the results:
# Execute the SQL query result_df <- sql(" SELECT user_id, date, status, COUNT(CASE WHEN status = 'moving' THEN 1 END) OVER ( PARTITION BY user_id ORDER BY UNIX_TIMESTAMP(date) * 1000 RANGE BETWEEN 5*3600*1000 PRECEDING AND CURRENT ROW ) AS rolling_5h_moving_count FROM device_status ORDER BY user_id, date ") # Show first 10 rows of results showDF(result_df, n = 10)
Key Details to Note
- Partitioning:
PARTITION BY user_idensures the rolling window is calculated separately for each user (no cross-user mixing). - Timestamp Conversion: We use
UNIX_TIMESTAMP(date) * 1000to convert the POSIX timestamp to milliseconds—Spark'sRANGEwindow requires a numeric sort key to define time ranges. - Window Range:
5*3600*1000is 5 hours in milliseconds. Adjust this number if you need a different window size (e.g., 2 hours would be2*3600*1000). - Filtered Count: The
CASE WHENinsideCOUNTensures we only count records where the status is'moving'. If you want to count all records in the window, replace that withCOUNT(*) OVER (...).
If your date column is stored as a string instead of a timestamp, modify the UNIX_TIMESTAMP call to match your date format, e.g.: UNIX_TIMESTAMP(date, 'yyyy-MM-dd HH:mm:ss') * 1000.
内容的提问来源于stack exchange,提问作者nate

