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

如何在PySpark中无需使用toPandas()为DataFrame记录分配统计频率,并将Python风格的编码与归一化预处理转为纯Spark实现?

Pure PySpark Implementation of Your Preprocessing Workflow

Hey there! Since you're new to PySpark and want to replicate your Pandas/Scikit-learn preprocessing without switching to a Pandas DataFrame, let's walk through the optimal pure-Spark approach step by step. We'll cover adding frequency counts, encoding categorical columns, and normalizing numeric columns—all while keeping everything in Spark's distributed framework.


1. Adding the Frequency Column

First, we need to calculate the relative frequency of each unique row (just like your add_freq_to_features function). Instead of Pandas-style groupby + merge, we'll use Spark's built-in aggregation to get this done efficiently for large datasets.

Approach:

  • Calculate the total number of records in your DataFrame.
  • Group by all columns to count occurrences of each unique row.
  • Compute relative frequency by dividing each group's count by the total record count.
  • Join this frequency data back to the original DataFrame to retain all rows.
from pyspark.sql import functions as F

# Step 1: Get total number of records
total_count = sdf.count()

# Step 2: Calculate group counts and relative frequencies
freq_df = sdf.groupBy(list(sdf.columns)) \
             .agg(F.count("*").alias("count")) \
             .withColumn("Freq", F.col("count") / total_count) \
             .drop("count")

# Step 3: Join frequency data back to original DataFrame
sdf_with_freq = sdf.join(freq_df, on=list(sdf.columns), how="left")

This will give you exactly the DataFrame you showed in your expected output, with the Freq column added.


2. Encoding Categorical Columns (A and C)

For categorical columns, Spark's StringIndexer is the direct equivalent of Scikit-learn's LabelEncoder. Unlike Pandas, we handle each categorical column individually (since StringIndexer operates on single columns).

Approach:

  • Use StringIndexer to convert categorical strings in columns A and C to numeric indices.
  • Fit the indexer on your data and transform it to get the encoded columns.
from pyspark.ml.feature import StringIndexer

# Encode column A
indexer_a = StringIndexer(inputCol="A", outputCol="A_encoded")
model_a = indexer_a.fit(sdf_with_freq)
sdf_encoded = model_a.transform(sdf_with_freq)

# Encode column C
indexer_c = StringIndexer(inputCol="C", outputCol="C_encoded")
model_c = indexer_c.fit(sdf_encoded)
sdf_encoded = model_c.transform(sdf_encoded)

Pro tip: If you expect unseen categories in future data, add handleInvalid="keep" to the StringIndexer constructor to avoid runtime errors.


3. Normalizing Numeric Columns

Spark's MinMaxScaler works on vector columns, so we'll first assemble all columns we want to normalize into a single vector, apply the scaler, then split the vector back into individual normalized columns. First, let's convert the boolean column D to numeric (0/1) since scalers only work on numeric data:

# Convert boolean column D to numeric (0/1)
sdf_encoded = sdf_encoded.withColumn("D_numeric", F.when(F.col("D") == True, 1).otherwise(0))

Now let's assemble and normalize:

from pyspark.ml.feature import VectorAssembler, MinMaxScaler

# List of columns to normalize (encoded categoricals, numeric cols, boolean-to-numeric, and Freq)
cols_to_normalize = ["A_encoded", "B", "C_encoded", "D_numeric", "E", "Freq"]

# Step 1: Assemble columns into a single vector
assembler = VectorAssembler(inputCols=cols_to_normalize, outputCol="features")
sdf_assembled = assembler.transform(sdf_encoded)

# Step 2: Apply MinMaxScaler
scaler = MinMaxScaler(inputCol="features", outputCol="normalized_features")
scaler_model = scaler.fit(sdf_assembled)
sdf_scaled = scaler_model.transform(sdf_assembled)

# Step 3: Split the normalized vector back into individual columns
for idx, col_name in enumerate(cols_to_normalize):
    sdf_scaled = sdf_scaled.withColumn(f"{col_name}_normalized", F.col("normalized_features")[idx])

# Optional: Drop intermediate columns if you don't need them
final_sdf = sdf_scaled.drop("features", "normalized_features", "A_encoded", "C_encoded", "D_numeric")

The final_sdf will have all your original columns plus their normalized versions (e.g., A_encoded_normalized, B_normalized) and the original Freq column (or its normalized version, depending on your needs).


Full Combined Code

Here's the complete code putting all steps together:

from pyspark.sql import functions as F
from pyspark.ml.feature import StringIndexer, VectorAssembler, MinMaxScaler

# ----------------------
# Step 1: Add Frequency Column
# ----------------------
total_count = sdf.count()
freq_df = sdf.groupBy(list(sdf.columns)) \
             .agg(F.count("*").alias("count")) \
             .withColumn("Freq", F.col("count") / total_count) \
             .drop("count")
sdf_with_freq = sdf.join(freq_df, on=list(sdf.columns), how="left")

# ----------------------
# Step 2: Encode Categorical Columns
# ----------------------
# Encode column A
indexer_a = StringIndexer(inputCol="A", outputCol="A_encoded")
model_a = indexer_a.fit(sdf_with_freq)
sdf_encoded = model_a.transform(sdf_with_freq)

# Encode column C
indexer_c = StringIndexer(inputCol="C", outputCol="C_encoded")
model_c = indexer_c.fit(sdf_encoded)
sdf_encoded = model_c.transform(sdf_encoded)

# Convert boolean D to numeric
sdf_encoded = sdf_encoded.withColumn("D_numeric", F.when(F.col("D") == True, 1).otherwise(0))

# ----------------------
# Step 3: Normalize Columns
# ----------------------
cols_to_normalize = ["A_encoded", "B", "C_encoded", "D_numeric", "E", "Freq"]

# Assemble columns into vector
assembler = VectorAssembler(inputCols=cols_to_normalize, outputCol="features")
sdf_assembled = assembler.transform(sdf_encoded)

# Scale the vector
scaler = MinMaxScaler(inputCol="features", outputCol="normalized_features")
scaler_model = scaler.fit(sdf_assembled)
sdf_scaled = scaler_model.transform(sdf_assembled)

# Split normalized vector back to columns
for idx, col_name in enumerate(cols_to_normalize):
    sdf_scaled = sdf_scaled.withColumn(f"{col_name}_normalized", F.col("normalized_features")[idx])

# Clean up intermediate columns
final_sdf = sdf_scaled.drop("features", "normalized_features", "A_encoded", "C_encoded", "D_numeric")

# Show the result
final_sdf.show()

Key Notes:

  • Efficiency: All operations are distributed, so they'll work on large datasets that can't fit in memory (unlike Pandas).
  • Flexibility: Adjust the cols_to_normalize list to include/exclude columns as needed.
  • Categorical Handling: StringIndexer follows the same logic as LabelEncoder—mapping the most frequent category to 0, next to 1, etc.

内容的提问来源于stack exchange,提问作者Mario

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 13:27:32