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

如何使用PySpark统计时间范围内事件每分钟的出现次数

PySpark: Count Events Per Minute Within a Specified Time Range

Got it, let's walk through how to solve this problem—counting the number of tweets per minute during the game's time window (from Game start minute to Game end minute), including minutes with zero events so we cover the full 130-minute range.

Step 1: Load and Prepare the Sample Data

First, let's turn your sample data into a PySpark DataFrame. I’ll assume your time values are integers representing consistent minute-based units (like epoch minutes):

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, lit, sequence, explode, count

# Initialize Spark session
spark = SparkSession.builder.appName("TweetMinuteCount").getOrCreate()

# Sample data mapped to columns
data = [
    (1001, 145678, 145600, 145730),
    (1002, 145678, 145600, 145730),
    (1005, 145680, 145600, 145730),
    (12278, 145687, 145600, 145730),
    (765558, 145688, 145600, 145730),
    (724323, 145689, 145600, 145730),
    (875857, 145688, 145600, 145730),
    (79375, 145685, 145600, 145730),
    (84666, 145686, 145600, 145730),
    (335556, 145687, 145600, 145730),
    (29990, 145688, 145600, 145730),
    (56, 145689, 145600, 145730),
    (968867, 145690, 145600, 145730),
    (8452, 145691, 145600, 145730),
    (1334, 145679, 145600, 145730)
]

# Create DataFrame with column names
df = spark.createDataFrame(data, ["tweet_id", "tweet_created_minute", "game_start_minute", "game_end_minute"])

Step 2: Generate the Full 130-Minute Time Range

To ensure we don’t miss minutes with zero tweets, we’ll create a DataFrame that includes every minute from the game’s start to end. Your sample uses a single consistent time range, so we can extract that once:

# Extract the unique game time range from the data
game_time = df.select("game_start_minute", "game_end_minute").distinct().collect()[0]
start_min = game_time.game_start_minute
end_min = game_time.game_end_minute

# Generate a sequence of all minutes in the range and explode into rows
full_minutes_df = spark.range(1) \
    .select(sequence(lit(start_min), lit(end_min)).alias("minutes")) \
    .select(explode(col("minutes")).alias("minute"))

Step 3: Count Tweets per Minute

Next, we filter tweets that fall within the game window (a safe practice even if all your sample data qualifies) and count how many occur each minute:

# Filter valid tweets and count per minute
tweet_counts_df = df.filter(col("tweet_created_minute").between(start_min, end_min)) \
    .groupBy("tweet_created_minute") \
    .agg(count("tweet_id").alias("event_count")) \
    .withColumnRenamed("tweet_created_minute", "minute")

Step 4: Combine and Fill Zero for Empty Minutes

Finally, do a left join between the full minute list and our tweet counts to include every minute in the 130-minute window. We’ll fill 0 for minutes with no tweets:

# Left join to retain all minutes, fill 0 where no tweets exist
final_result_df = full_minutes_df.join(tweet_counts_df, on="minute", how="left") \
    .fillna(0, subset=["event_count"]) \
    .orderBy("minute")

# View the result (you can also write to storage with .write)
final_result_df.show(15)  # Show first 15 rows as an example

Quick Notes for Edge Cases:

  • Variable Game Times: If different rows have unique start/end times, you’d need to generate time ranges per game and adjust grouping logic accordingly.
  • Timestamp Values: If your time data is in timestamp format instead of integers, use date_trunc("minute", col("tweet_created_time")) to get minute-level granularity.
  • Performance: For large datasets, partition your DataFrame by time ranges to optimize joins and grouping operations.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 09:39:19