Spark实现CSV中匹配数组列的对应式拆分问询
Got it, let's break down how to split your two array columns into matching pairs without generating a Cartesian product. The key here is using Spark's arrays_zip function to tie together corresponding elements from each array, then exploding the result to get individual rows.
Step 1: Fix Data Types (Critical!)
First, when reading CSV files, Spark doesn't automatically recognize string representations of arrays (like "[11,23]") as actual array types. We need to convert those string columns into proper arrays first:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.types._ // Read the CSV (add .option("header", "true") if your file has column headers) val table = spark.read.csv("my_data.csv") // Convert string columns to actual arrays val dfWithArrays = table // Clean up the string: remove brackets, split by commas, then cast to integer array for ids .withColumn("ids", split(regexp_replace(col("_c0"), "\\[|\\]", ""), ",").cast(ArrayType(IntegerType))) // Same for avg: cast to double array .withColumn("avg", split(regexp_replace(col("_c1"), "\\[|\\]", ""), ",").cast(ArrayType(DoubleType))) // Drop the original string columns .drop("_c0", "_c1")
Step 2: Zip and Explode Corresponding Elements
Now use arrays_zip to combine the two arrays into an array of structs (each struct holds one id-avg pair), then explode that array to get each pair as a separate row:
val finalResult = dfWithArrays // Zip the two arrays into an array of (id, avg) structs .withColumn("zipped_pairs", arrays_zip(col("ids"), col("avg"))) // Explode the zipped array to create individual rows for each pair .select(explode(col("zipped_pairs")).alias("pair")) // Extract the id and avg from the struct and rename columns .select(col("pair.ids").alias("id"), col("pair.avg").alias("avg"))
Expected Output
For your sample data, this will produce:
+---+--------+ |id |avg | +---+--------+ |11 |0.368633| |23 |0.750615| +---+--------+
Notes
- If your CSV includes headers (like
ids,avg), add.option("header", "true")when reading the file. Then you can reference columns by their names directly instead of_c0and_c1. - This approach ensures you only get matching pairs (no Cartesian product) because
arrays_zipstrictly pairs elements at the same index from each array.
内容的提问来源于stack exchange,提问作者Iza Pziza

