apache / apache/datafusion

Create memory table with target partitions

Open
#12,905 4 comments 0 reactions 1 assignee Claimed by @jayzhan211 View on GitHub
bug
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

### Describe the bug

When we create table inside slt, it is created with MemoryExec (Memory table). But it seems like the partitions is always one. If we want to create table with multiple partitions, we need to create multiple table and union with them.

### To Reproduce

```

statement ok
create table t1(a int, b varchar) as values
(1, 'a'),
(1, 'a'),
(1, 'a');

statement ok
create table t2(a int, b varchar) as values
(2, 'b'),
(2, 'b'),
(2, 'b');

query TI
With T as
(select * from t1 UNION ALL select * from t2)
select b, sum(DISTINCT a) from T group by b;
----
a 1
b 2

query TT
explain With T as
(select * from t1 UNION ALL select * from t2)
select b, sum(DISTINCT a) from T group by b;
----
logical_plan
01)Projection: t.b, sum(alias1) AS sum(DISTINCT t.a)
02)--Aggregate: groupBy=[[t.b]], aggr=[[sum(alias1)]]
03)----Aggregate: groupBy=[[t.b, CAST(t.a AS Int64) AS alias1]], aggr=[[]]
04)------SubqueryAlias: t
05)--------Union
06)----------TableScan: t1 projection=[a, b]
07)----------TableScan: t2 projection=[a, b]
physical_plan
01)ProjectionExec: expr=[b@0 as b, sum(alias1)@1 as sum(DISTINCT t.a)]
02)--AggregateExec: mode=FinalPartitioned, gby=[b@0 as b], aggr=[sum(alias1)]
03)----CoalesceBatchesExec: target_batch_size=8192
04)------RepartitionExec: partitioning=Hash([b@0], 4), input_partitions=4
05)--------AggregateExec: mode=Partial, gby=[b@0 as b], aggr=[sum(alias1)]
06)----------AggregateExec: mode=FinalPartitioned, gby=[b@0 as b, alias1@1 as alias1], aggr=[]
07)------------CoalesceBatchesExec: target_batch_size=8192
08)--------------RepartitionExec: partitioning=Hash([b@0, alias1@1], 4), input_partitions=4
09)----------------RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=2
10)------------------AggregateExec: mode=Partial, gby=[b@1 as b, CAST(a@0 AS Int64) as alias1], aggr=[]
11)--------------------UnionExec
12)----------------------MemoryExec: partitions=1, partition_sizes=[1]
13)----------------------MemoryExec: partitions=1, partition_sizes=[1]

statement ok
set datafusion.execution.target_partitions = 2

statement ok
create table t(a int, b varchar) as values
(1, 'a'),
(1, 'a'),
(1, 'a'),
(2, 'b'),
(2, 'b'),
(2, 'b');

query TI
select b, sum(DISTINCT a) from T group by b;
----
b 2
a 1

query TT
explain select b, sum(DISTINCT a) from T group by b;
----
logical_plan
01)Projection: t.b, sum(alias1) AS sum(DISTINCT t.a)
02)--Aggregate: groupBy=[[t.b]], aggr=[[sum(alias1)]]
03)----Aggregate: groupBy=[[t.b, CAST(t.a AS Int64) AS alias1]], aggr=[[]]
04)------TableScan: t projection=[a, b]
physical_plan
01)ProjectionExec: expr=[b@0 as b, sum(alias1)@1 as sum(DISTINCT t.a)]
02)--AggregateExec: mode=FinalPartitioned, gby=[b@0 as b], aggr=[sum(alias1)]
03)----CoalesceBatchesExec: target_batch_size=8192
04)------RepartitionExec: partitioning=Hash([b@0], 2), input_partitions=2
05)--------AggregateExec: mode=Partial, gby=[b@0 as b], aggr=[sum(alias1)]
06)----------AggregateExec: mode=FinalPartitioned, gby=[b@0 as b, alias1@1 as alias1], aggr=[]
07)------------CoalesceBatchesExec: target_batch_size=8192
08)--------------RepartitionExec: partitioning=Hash([b@0, alias1@1], 2), input_partitions=2
09)----------------RepartitionExec: partitioning=RoundRobinBatch(2), input_partitions=1
10)------------------AggregateExec: mode=Partial, gby=[b@1 as b, CAST(a@0 AS Int64) as alias1], aggr=[]
11)--------------------MemoryExec: partitions=1, partition_sizes=[1] <- It is still 1

```

### Expected behavior

I hope we can create arbitrary partitions with setting instead of creating the test with union.

### Additional context

I think batch size doesn't not change the MemoryExec too, but I'm not sure about the intended behaviour

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.