join two date partitioned parquet hive tables on specific dates

Viewed 39

I have two hive parquet tables which are partitioned based on dates, both the tables are having several millions of records, I need to join the two tables based on hundreds of specific dates and select around 150 columns. When I tried the below SQL, the query is running for ever. I'm running this SQL in pyspark, Is there any other way to optimize it?

SQL1:

select table_a.a_col1, table_a.a_col2, table_b.b_col1, table_b.b_col2,..., table_a.a_col150, table_b.b_col150
from table_a a join table_b b
on a.col1 = b.col2
and a.date_col = b.date_col
where a.date_col in ('2022-01-01','2022-01-02','2022-01-03','2022-01-04')
and b.date_col in ('2022-01-01','2022-01-02','2022-01-03','2022-01-04')

SQL2:

select table_a.a_col1, table_a.a_col2, table_b.b_col1, table_b.b_col2,..., table_a.a_col150, table_b.b_col150
from table_a a join table_b b
on a.col1 = b.col2
where a.date_col in ('2022-01-01','2022-01-02','2022-01-03','2022-01-04')
and b.date_col in ('2022-01-01','2022-01-02','2022-01-03','2022-01-04')
1 Answers

You can store all dates into a lkp table and the join it to table_a, table_b.

step 1 - first create a lookup table - create table lkp_date_filter (dt timestamp);
step 2 - insert filter dates into it - insert into lkp_date_filter values('2022-01-04')
step 3 - join it in your main query and remove IN claues.

select table_a.a_col1, table_a.a_col2, table_b.b_col1, table_b.b_col2,..., table_a.a_col150, table_b.b_col150
from table_a a 
join table_b b on a.col1 = b.col2
join lkp_date_filter on a.date_col =lkp.dt and b.date_col =lkp.dt

Step 3 will avoid the expensive IN clause and make SQL fast. Step2 will give you flexibility to change the filter values as per your need. You can partition table a and b on date_col to make SQL faster.

Related