sqlContext.sql("select id, city, value from A").show()
+----+-------+-------+
| id | city | value |
+----+-------+-------+
| id1| city1 | 4.0 |
+----+-------+-------+
sqlContext.sql("select id, city, rate from B").show()
+----+-------+-------+
| id | city | rate |
+----+-------+-------+
| id1| city1 | 1.0 |
+----+-------+-------+
sqlContext.sql("""
select
a.id as a_id
, b.id as b_id
, a.city as a_city
, b.city as b_city
, a.value
, b.rate
from A a left join B b
on a.id = b.id and a.city= b.city
""").show()
+------+------+---=----+--------+-------+------+
| a_id | b_id | a_city | b_city | value | rate |
+------+------+---=----+--------+-------+------+
| id1 | id1 | city1 | city1 | 4.0 | 0.0 |
+------+------+---=----+--------+-------+------+
b.rate should be 1.0. but it is 0.0
it is working in spark 2.3 but not in spark 3.1
left outer join is not working too.
Looks like a version problem, I looked at the release note, but couldn't find a suitable explanation.