Spark DataFrame基于旧列新增Voltage与Current列的Scala实现咨询
Solution for Adding Voltage and Current Columns in Spark Scala
Here's the Scala Spark code to achieve your requirement, with clear explanations for each step:
First, ensure you have the necessary imports for Spark's core and function libraries:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.{when, col}
Initialize your SparkSession (skip this block if you already have an active session running):
val spark = SparkSession.builder() .appName("VoltageCurrentColumnAddition") .master("local[*]") // Remove this line when deploying to a cluster .getOrCreate()
Create your input DataFrame using the sample data provided:
val inputDF = spark.createDataFrame(Seq( ("inv_1_c", 5), ("inv_1_v", 8), ("inv_2_c", 9) )).toDF("key", "value")
Add the Voltage and Current columns using conditional logic:
val resultDF = inputDF .withColumn("Voltage", when(col("key").endsWith("_v"), col("value")).otherwise(0)) .withColumn("Current", when(col("key").endsWith("_c"), col("value")).otherwise(0))
To view the final output, run:
resultDF.show()
Key Explanations:
when(condition, valueIfTrue).otherwise(valueIfFalse): This Spark function implements conditional logic. ForVoltage, we check if thekeyends with "_v" — if yes, we take the correspondingvalue, else we use 0. The same logic applies toCurrentwith the "_c" suffix.col("key").endsWith("_v"): Uses Spark's built-in string function to verify the suffix of thekeycolumn.withColumn: Adds a new column to the DataFrame (or replaces an existing column with the same name).
Expected Output:
+--------+-----+-------+-------+ | key|value|Voltage|Current| +--------+-----+-------+-------+ |inv_1_c | 5| 0| 5| |inv_1_v | 8| 8| 0| |inv_2_c | 9| 0| 9| +--------+-----+-------+-------+
内容的提问来源于stack exchange,提问作者Soumyadip Ghosh
相关产品推荐
相关产品推荐

