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

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

  1. Mismatched Regular Expressions
    Your regex patterns don't align with your actual log format:

    • For host, your regex r'^([^\s]+\s)' captures the hostname plus an extra space (since it includes \s at 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.
  2. No Time Window Logic
    You're doing a simple groupBy('host').sum(...), which gives you the total ByteSize per hostname across all time—not per hourly tumbling window, which is what you need.

  3. Not Using the Aggregated Result
    You run split_df.groupby(...).sum(...) but don't assign this to a variable, then call split_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 Timestamp type to handle time windows, so we converted the raw dd:hh:mm:ss string 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 the groupBy tells 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_df and 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 06:33:12