PySpark:计算忽略零值与空值的行最小值并生成新列v5
v5 in PySpark DataFrame Based on Key-Value Pairs Got it, let's break down how to create the v5 column you need in PySpark. The core idea is to first filter out invalid values (zeros and nulls) from v1 and v2, then pick the smallest valid key and map it to its corresponding value (v3 for v1, v4 for v2).
Step 1: Set Up Example DataFrame
First, let's create a sample DataFrame to test our logic with different edge cases (like nulls, zeros, valid values for both keys):
from pyspark.sql import SparkSession from pyspark.sql import functions as F # Initialize Spark session spark = SparkSession.builder.appName("CreateV5Column").getOrCreate() # Sample data covering various scenarios data = [ (1, 3, "a", "c"), # Both v1 and v2 valid; v1 is smaller (0, 2, "b", "d"), # v1 is zero (invalid), only v2 valid (None, 5, "e", "f"),# v1 is null (invalid), only v2 valid (4, 0, "g", "h"), # v2 is zero (invalid), only v1 valid (None, None, "i", "j"), # Both keys invalid (2, 1, "k", "l") # Both valid; v2 is smaller ] df = spark.createDataFrame(data, ["v1", "v2", "v3", "v4"]) df.show()
Step 2: Implement the Logic with Conditional Statements
This approach uses nested when() clauses to handle each scenario explicitly:
df_with_v5 = df.withColumn( "v5", # Case 1: Only v1 is valid (non-null, non-zero) F.when( (F.col("v1").isNotNull() & (F.col("v1") != 0)) & ~(F.col("v2").isNotNull() & (F.col("v2") != 0)), F.col("v3") ) # Case 2: Only v2 is valid .when( ~(F.col("v1").isNotNull() & (F.col("v1") != 0)) & (F.col("v2").isNotNull() & (F.col("v2") != 0)), F.col("v4") ) # Case 3: Both keys are valid; pick the smaller key's value .when( (F.col("v1").isNotNull() & (F.col("v1") != 0)) & (F.col("v2").isNotNull() & (F.col("v2") != 0)), F.when(F.col("v1") < F.col("v2"), F.col("v3")).otherwise(F.col("v4")) ) # Case 4: Neither key is valid; set to null .otherwise(None) ) df_with_v5.show()
Expected Output:
+----+----+---+---+----+ | v1 | v2 |v3 |v4 | v5 | +----+----+---+---+----+ | 1 | 3 |a |c | a | | 0 | 2 |b |d | d | |null| 5 |e |f | f | | 4 | 0 |g |h | g | |null|null|i |j |null| | 2 | 1 |k |l | l | +----+----+---+---+----+
Alternative Approach: Using Arrays and Structs
If you anticipate adding more key-value pairs later, this method is more scalable. We'll create an array of valid key-value structs, filter out invalid entries, then pick the minimum key's value:
df_with_v5_alt = df.withColumn( # Create an array of structs (key, value) for valid entries "valid_pairs", F.array( F.when(F.col("v1").isNotNull() & (F.col("v1") != 0), F.struct(F.col("v1").alias("key"), F.col("v3").alias("value"))), F.when(F.col("v2").isNotNull() & (F.col("v2") != 0), F.struct(F.col("v2").alias("key"), F.col("v4").alias("value"))) ) ) # Filter out null entries from the array .withColumn( "valid_pairs_filtered", F.filter("valid_pairs", lambda x: x.isNotNull()) ) # Get the struct with the smallest key .withColumn( "min_pair", F.array_min("valid_pairs_filtered") ) # Extract the value from the min struct to get v5 .withColumn( "v5", F.col("min_pair.value") ) # Clean up intermediate columns .drop("valid_pairs", "valid_pairs_filtered", "min_pair") df_with_v5_alt.show()
This will produce the exact same output as the first method, but it's easier to extend if you add more columns like v6/v7 later—just add another when() clause to the valid_pairs array.
内容的提问来源于stack exchange,提问作者Mia21

