PySpark中无法基于1小时Tumbling Window进行Group By分组问题
Hey there! Let's walk through what's off with your current code and fix it up to meet your requirement—grouping by hostname and summing ByteSize over hourly tumbling windows. I'll break down the issues first, then show you the corrected code with explanations.
Key Issues in Your Current Code
Mismatched Regular Expressions
Your regex patterns don't align with your actual log format:- For
host, your regexr'^([^\s]+\s)'captures the hostname plus an extra space (since it includes\sat the end). - For
timestamp, you're trying to match a format like[dd/MMM/yyyy:hh:mm:ss -xxxx], but your logs use[dd:hh:mm:ss]—so this regex returns empty values entirely.
Without correct timestamp extraction, you can't build time-based windows at all.
- For
No Time Window Logic
You're doing a simplegroupBy('host').sum(...), which gives you the total ByteSize per hostname across all time—not per hourly tumbling window, which is what you need.Not Using the Aggregated Result
You runsplit_df.groupby(...).sum(...)but don't assign this to a variable, then callsplit_df.show()—so you're only showing the raw parsed data, not your aggregated stats.
Corrected Code
Let's fix each issue step by step:
import findspark findspark.init() from pyspark.sql import SparkSession from pyspark.sql.functions import regexp_extract, to_timestamp, concat_ws, split, window # Initialize SparkSession (modern Spark uses this instead of SQLContext) spark = SparkSession.builder.appName("LogAnalysis").getOrCreate() # Read the log file base_df = spark.read.text("C:/Users/Documents/log.txt") # Correct regex extraction to match your log format split_df = base_df.select( # Extract hostname: match all characters from start to first space regexp_extract('value', r'^([^\s]+)', 1).alias('host'), # Extract raw timestamp from []: matches dd:hh:mm:ss regexp_extract('value', r'\[(\d{2}:\d{2}:\d{2}:\d{2})]', 1).alias('raw_timestamp'), # Extract ByteSize: match the final number in the line regexp_extract('value', r'^.*\s+(\d+)$', 1).cast('integer').alias('content_size') ) # Convert raw timestamp to Spark's Timestamp type (required for window operations) # Note: Your logs only have day:hour:minute:ss, so we'll add a dummy date to make it valid formatted_df = split_df.withColumn( 'timestamp', to_timestamp( concat_ws( ' ', # Combine dummy date with the day from raw timestamp concat_ws('-', '2024', '01', split('raw_timestamp', ':')[0]), # Combine hour:minute:second concat_ws(':', split('raw_timestamp', ':')[1], split('raw_timestamp', ':')[2], split('raw_timestamp', ':')[3]) ), 'yyyy-MM-dd HH:mm:ss' ) ) # Group by hostname AND hourly tumbling window, then sum content_size result_df = formatted_df.groupBy( 'host', window('timestamp', '1 hour') # Defines a 1-hour tumbling window ).sum('content_size').withColumnRenamed('sum(content_size)', 'total_bytes') # Show the aggregated results result_df.show(truncate=False)
What Changed?
- Fixed Regex: Now correctly captures hostname (no extra space) and the timestamp from your log format.
- Timestamp Conversion: Spark needs a proper
Timestamptype to handle time windows, so we converted the rawdd:hh:mm:ssstring to a full datetime (using a dummy date—you can adjust this if your logs have actual dates). - Hourly Tumbling Window: Using
window('timestamp', '1 hour')in thegroupBytells Spark to split data into non-overlapping 1-hour chunks, then sum the ByteSize per hostname per window. - Used Aggregated Result: We assigned the grouped/summed data to
result_dfand show that, instead of the raw parsed data.
If your logs actually include full dates (not just day:hour:minute:ss), just adjust the to_timestamp format string to match your actual timestamp pattern—this code is easy to tweak for that!
内容的提问来源于stack exchange,提问作者user2961127

