We are moving a database to a delta lake, and have copied the tables to delta tables in azure delta lake gen2.
The users access the data through a serverless sql pool in synapse.
We have run into a showstopper.
When doing these queries table like this:
select top 1000
col1
,col2
,...
,col50
from tablea
select top 1000
col1
,col2
,...
,col60
from tableb
each will take about 5 seconds at most. the data use in each when looking at the look might ba about 100 MB
when doing this query
select top 1000 * from
(
select
col1
,col2
,...
,col50
from tablea
union all
select top 1000
col1
,col2
,...
,col50
from tableb
) a
it takes minutes, and the log shows a memory use of hundreds of gigabytes
Naively I would think that the optimizer would just return the first 1000 rows from the first table in the union all, but for some reason it parses all of the data before returning the first 1000 rows.
The only solution we can see is to create a new set of tables in the delta lake, with the tables unioned, however this seems like a waste of space.
We have seperate tables, because the dimensionality is slightly different, and the measures are wholly different.
In perhaps 80% of the time the users will only query rows from one of the tables.
Is there a way to optimize union all in serverless queries?