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

如何从已读取HDFS数据的PySpark DataFrame提取数值列并计算Z-score

Filter Numeric Columns and Calculate Z-Scores in PySpark

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:44:54