Scala实现基于两表关联字段新增存在性标记列
You can build your desired table3 using either Spark SQL or the DataFrame API—both are straightforward given your existing DataFrames df (table1) and df1 (table2). Here are step-by-step implementations for both methods:
Approach 1: Using Spark SQL Directly
This method leverages a left join and a CASE statement to determine the IND value, which is easy to read and aligns with standard SQL practices:
// Run the SQL query to generate table3 and create a temporary view hc.sql(""" SELECT t1.id, t1.name, CASE WHEN t2.id IS NOT NULL THEN 'Y' ELSE 'N' END AS IND FROM table1 t1 LEFT JOIN table2 t2 ON t1.id = t2.id """).createOrReplaceTempView("table3") // Optional: Save as a permanent table (overwrites if exists) // hc.sql("CREATE OR REPLACE TABLE table3 AS SELECT ...")
The left join ensures all rows from table1 are preserved. For rows where id exists in table2, t2.id will have a value (so IND = 'Y'), and for non-matching rows, t2.id will be null (so IND = 'N').
Approach 2: Using DataFrame API
If you prefer working with DataFrame operations directly, use a left outer join and Spark's when function to derive the IND column:
First, import the required Spark functions:
import org.apache.spark.sql.functions.{when, col}
Then, perform the join and construct the final DataFrame:
// Left join table1 and table2 on the id column val joinedDf = df.join(df1, df("id") === df1("id"), "left_outer") // Add the IND column and select only the necessary fields val table3Df = joinedDf.select( df("id").alias("id"), df("name").alias("name"), when(df1("id").isNotNull, "Y").otherwise("N").alias("IND") ) // Create a temporary view for table3 table3Df.createOrReplaceTempView("table3") // Optional: Save as a permanent table // table3Df.write.mode("overwrite").saveAsTable("table3")
Both approaches will generate exactly the table3 output you specified, with the correct IND values for each row.
内容的提问来源于stack exchange,提问作者Sudheer Nulu

