如何使用PySpark统计时间范围内事件每分钟的出现次数
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

