基于PySpark统计当前时间N天内各ID的事件发生次数
Alright, let's solve this problem where we need to count how many times each ID has had events in the last N days relative to the current time. Here's a step-by-step approach using PySpark:
Step 1: Set Up Dependencies and Sample Data
First, let's import the necessary PySpark modules and create your sample DataFrame to work with:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, current_timestamp, count, expr from pyspark.sql.types import TimestampType # Initialize Spark session spark = SparkSession.builder.appName("EventCountByID").getOrCreate() # Sample data data = [ ("533ladk203ldpwk", "2018-03-28 17:52:04"), ("516dlksw9823adp", "2018-03-26 12:58:04"), ("516dlksw9823adp", "2018-01-24 07:52:16"), ("533ladk203ldpwk", "2018-03-18 03:23:11"), ("533ladk203ldpwk", "2018-03-14 08:30:13") ] # Create DataFrame and convert Occurrence to TimestampType df = spark.createDataFrame(data, ["id", "Occurrence"]) df = df.withColumn("Occurrence", col("Occurrence").cast(TimestampType()))
Step 2: Define the Time Window and Filter Data
Next, we'll set our N-day window (let's use N=10 as an example) and calculate the cutoff time (current time minus N days). Then we'll filter the DataFrame to only include events that happened after this cutoff:
# Set the number of days for the window N = 10 # Calculate cutoff time: current time minus N days cutoff_time = current_timestamp() - expr(f"interval {N} days") # Filter events within the last N days and count per ID event_counts = df.filter(col("Occurrence") >= cutoff_time) \ .groupBy("id") \ .agg(count("Occurrence").alias(f"event_count_last_{N}_days"))
Step 3: View the Results
Finally, let's show the output. Using our sample data and assuming the current time is 2018-03-29 00:00:00, the cutoff would be 2018-03-19 00:00:00. This means:
533ladk203ldpwkonly has 1 event (2018-03-28) in the last 10 days516dlksw9823adphas 1 event (2018-03-26) in the last 10 days- The older events (2018-01-24, 2018-03-18, 2018-03-14) are excluded
To see the results:
event_counts.show()
Expected output:
+---------------+-----------------------+ | id|event_count_last_10_days| +---------------+-----------------------+ |533ladk203ldpwk| 1| |516dlksw9823adp| 1| +---------------+-----------------------+
Quick Notes
- If you need a fixed cutoff time instead of relative to current time, replace
current_timestamp() - expr(f"interval {N} days")with a hardcoded timestamp (e.g.,to_timestamp("2018-03-19 00:00:00")). - Double-check that your
Occurrencecolumn is properly cast toTimestampType—this is critical for accurate date comparisons.
内容的提问来源于stack exchange,提问作者Ashley O

