Spark SQL tests all types of join in JoinType for ease of understanding
Prepare test data
Trade
Order number seller buyer city
1 A Xiao Wang Beijing 2 B Xiao Li Tianjin 3 A Xiao Liu Beijing
Order
Order number buyer's commodity name price delivery time
1 Xiao Wang TV 12 2015-08-01 09 purl 08RU 311 Xiao Wang refrigerator 24 2015-08-01 09RV 08Rd 142Xiao Li air conditioner 12 2015-09-02 09:01:31
Note: all are divided by\ t
Create DF def main (args: Array [String]): Unit = {val spark=SparkSession.builder () .appName ("JoinDemo") .master ("local [2]") .getOrCreate () import spark.implicits._ val order=spark.sparkContext.textFile ("order.data"). Map (_ .split ("\ t")) .map (x = > Order (x (0), x (1), x (2), x (3)) X (4)) .toDF () val trade=spark.sparkContext.textFile ("trade.data") .map (_ .split ("\ t")) .map (x = > Trade (x (0), x (1), x (2)) X (3)). ToDF () order.show () / / +-- + / | o_id | buyer | p_name | price | date | / / +-- + / / | 1 | Xiao Wang | TV | 12 | 2015-08-01 | / / | 1 | Xiao Wang | refrigerator | 24 | 2015-08-01 | / / | 2 | Xiao Li | Air conditioning | 12 | 2015-09-02 | / / +-+ trade.show () / / +-- +-+ / / | o_id | seller | buyer | buyer_city | / / +-- + / / | 1 | A | Xiao Wang | Beijing | / / | 2 | B | Xiao Li | Tianjin | / / | 3 | A | Xiao Liu | Beijing | / / + -+} case class Student (id:String Name:String,phoneNum:String,email:String) case class Order (obliquidPregnameGregory string gimmick) case class Trade (obliquidPregnameParticipate StringGetWord (date) string) case class Trade (obsolidPregnameParticipationStringparentin buyerGlyString.buyerGlyString.buyercase class Trade) string type defaults to `inner`.String. Must be one of the following types: `inner`, `cross`, `outer`, `full`, `full_ outer`, `left`, `left_ outer`, `right`, `right_ outer`, `left_ semi`, `left_ anti`.1, unspecified and inner
Do not specify
Trade.join (order Trade ("o_id") = = order ("o_id") .show+----+----+ | o_id | seller | buyer | buyer_city | buyer | p_name | price | date | + -- +-- + | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | TV | 12 | 2015-08-01 | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | 24 | 2015-08-01 | refrigerator | 2 | B | Xiao Li | Tianjin | 2 | Xiao Li | Air conditioning | 12 | 2015-09-02 | +-+
Specify inner
Scala > trade.join (order,trade ("o_id") = order ("o_id") "inner"). Show+----+----+ | o_id | seller | buyer | buyer_city | o_id | buyer | p_name | price | date | +- -- + | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | TV | 12 | 2015-08-01 | | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | refrigerator | 24 | 2015-08-01 | | 2 | B | Xiao Li | Tianjin | 2 | Xiao Li | Air conditioning | 12 | 2015-09-02 | + -+
Unspecified is the same as inner, which is to find the intersection of two Datarame.
2. Left and left outerscala > trade.join (order,trade ("o_id") = = order ("o_id") "left"). Show+----+----+ | o_id | seller | buyer | buyer_city | o_id | buyer | p_name | price | date | +- -- + | 3 | A | Xiao Liu | Beijing | null | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | TV | 12 | 2015-08-01 | | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | 24 | 2015-08-01 | | 2 | B | Xiao Li | Tianjin | 2 | Xiao Li | Air conditioning | 12 | 2015-09-02 | +-+ scala > trade.join (order) Trade ("o_id") = order ("o_id") "left_outer"). Show+----+----+ | o_id | seller | buyer | buyer_city | o_id | buyer | p_name | price | date | +- -+ | 3 | A | Xiao Liu | Beijing | null | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | TV | 12 | 2015-08-01 | | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | 24 | 2015-08-01 | | 2 | B | Xiao Li | | Tianjin | 2 | Xiao Li | Air conditioning | 12 | 2015-09-02 | +-+ |
Left join and left outer join are completely equivalent.
Right and right outerscala > trade.join (order,trade ("o_id") = = order ("o_id") "right_outer"). Show+----+----+ | o_id | seller | buyer | buyer_city | o_id | buyer | p_name | price | date | +- -+ | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | TV | 12 | 2015-08-01 | | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | refrigerator | 24 | 2015-08-01 | | 2 | B | Xiao Li | Tianjin | 2 | Xiao Li | Air conditioning | 12 | 2015-09 -02 | +-+ scala > trade.join (order Trade ("o_id") = order ("o_id") "right"). Show+----+----+ | o_id | seller | buyer | buyer_city | o_id | buyer | p_name | price | date | + +-+ | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | TV | 12 | 2015-08-01 | | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | refrigerator | 24 | 2015-08-01 | | 2 | B | Xiao Li | Tianjin | 2 | Xiao Li | Air conditioning | 12 | 2015-09-02 | + -+
Right and right outer are equivalent
Full_outerscala > trade.join (order,trade ("o_id") = order ("o_id") "full_outer"). Show+----+----+ | o_id | seller | buyer | buyer_city | o_id | buyer | p_name | price | date | +- -+ | 3 | A | Xiao Liu | Beijing | null | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | TV | 12 | 2015-08-01 | | 1 | A | Xiao Wang | Beijing | 1 | Xiao Wang | 24 | 2015-08-01 | | 2 | B | Xiao Li | | Tianjin | 2 | Xiao Li | Air conditioning | 12 | 2015-09-02 | +-+ |
Get the intersection of two Datarame
Left_semiscala > trade.join (order,trade ("o_id") = order ("o_id") "left_semi") .show+----+ | o_id | seller | buyer | buyer_city | +-- + | 1 | A | Xiao Wang | Beijing | | 2 | B | Xiao Li | Tianjin | +-+
Filter out the common parts of the two DF
Left_anticala > trade.join (order,trade ("o_id") = order ("o_id") "left_anti") .show+----+ | o_id | seller | buyer | buyer_city | +-+ | 3 | A | Xiao Liu | Beijing | +-+
Filter out the parts of DF2 that DF1 does not have