Union all performance on serverless sql in synapse is poor

Viewed 50

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?

0 Answers
Related