如何从已读取HDFS数据的PySpark DataFrame提取数值列并计算Z-score
Got it, let's break down how to tackle this task step by step. You've already loaded your JSON data from HDFS into a PySpark DataFrame, so we can jump right into the next parts.
Step 1: Identify Numeric Columns
First, we need to pick out all columns that hold numerical values (integers, floats, doubles) since Z-scores only make sense for quantitative data. Here are two reliable ways to do this:
from pyspark.sql.types import IntegerType, DoubleType, FloatType # Option 1: Simpler string matching with df.dtypes numeric_cols = [col_name for col_name, dtype in df.dtypes if dtype in ("int", "double", "float")] # Option 2: Precise type checking using the DataFrame schema (better for edge cases) numeric_cols = [ field.name for field in df.schema.fields if isinstance(field.dataType, (IntegerType, DoubleType, FloatType)) ]
Step 2: Calculate Mean and Standard Deviation
Next, we need the mean and standard deviation for each numeric column—these are the core values needed to compute Z-scores. We can grab all these stats in one pass using PySpark's agg function:
from pyspark.sql.functions import mean, stddev # Compute aggregate stats for all numeric columns stats_df = df.agg( *[mean(col).alias(f"{col}_mean") for col in numeric_cols], *[stddev(col).alias(f"{col}_stddev") for col in numeric_cols] ) # Convert stats to a dictionary for easy access stats_dict = stats_df.collect()[0].asDict()
Step 3: Compute Z-Scores for Each Numeric Column
Now we can add new columns to our DataFrame that hold the Z-score for each numeric column. The Z-score formula is:Z = (value - mean) / standard_deviation
We'll loop through each numeric column and apply this formula, with a check to avoid division by zero:
from pyspark.sql.functions import col, lit # Iterate over numeric columns to add Z-score columns for col_name in numeric_cols: mean_val = stats_dict[f"{col_name}_mean"] stddev_val = stats_dict[f"{col_name}_stddev"] # Handle columns with zero standard deviation (all values are identical) if stddev_val != 0: df = df.withColumn(f"{col_name}_zscore", (col(col_name) - mean_val) / stddev_val) else: df = df.withColumn(f"{col_name}_zscore", lit(0.0))
Alternative: Per-Group Z-Scores (Using Window Functions)
If you need to calculate Z-scores grouped by a categorical column (instead of across the entire dataset), use window functions instead:
from pyspark.sql.window import Window # Define your grouping window (replace "group_col" with your actual grouping column) window_spec = Window.partitionBy("group_col") # Calculate Z-scores for each group for col_name in numeric_cols: df = df.withColumn( f"{col_name}_zscore", (col(col_name) - mean(col_name).over(window_spec)) / stddev(col_name).over(window_spec) )
Verify Your Results
To check the output, run this to see original values alongside their Z-scores:
df.select(numeric_cols + [f"{col}_zscore" for col in numeric_cols]).show()
内容的提问来源于stack exchange,提问作者Praveen Mandadi

