如何在PySpark中无需使用toPandas()为DataFrame记录分配统计频率,并将Python风格的编码与归一化预处理转为纯Spark实现?
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
StringIndexerto 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_normalizelist to include/exclude columns as needed. - Categorical Handling:
StringIndexerfollows the same logic asLabelEncoder—mapping the most frequent category to 0, next to 1, etc.
内容的提问来源于stack exchange,提问作者Mario

