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

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_id ensures the rolling window is calculated separately for each user (no cross-user mixing).
  • Timestamp Conversion: We use UNIX_TIMESTAMP(date) * 1000 to convert the POSIX timestamp to milliseconds—Spark's RANGE window requires a numeric sort key to define time ranges.
  • Window Range: 5*3600*1000 is 5 hours in milliseconds. Adjust this number if you need a different window size (e.g., 2 hours would be 2*3600*1000).
  • Filtered Count: The CASE WHEN inside COUNT ensures we only count records where the status is 'moving'. If you want to count all records in the window, replace that with COUNT(*) 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 10:35:19