Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 25 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -171,7 +171,8 @@ so if you forget.

Comments (`-- ...` and `/* ... */`) can go anywhere. A syntax error reports the
line and column it was found at, and features efsql doesn't support (`OR`,
`NOT`, `<>`, joins, functions, `GROUP BY`) are rejected by name.
`NOT`, `<>`, joins, `HAVING`, functions other than aggregates) are rejected by
name.

### Storage IDs

Expand Down Expand Up @@ -303,6 +304,29 @@ Typed literals work anywhere a value does, including `IN` and `BETWEEN`:
select col_a from tenant_id.table_name where inserted_at between '2024-01-01'::timestamp and '2025-01-01'::timestamp;
```

### Group and aggregate

```sql
-- one row for the whole table
select count(*) from tenant_id.table_name;
select min(inserted_at), max(inserted_at) from tenant_id.table_name where status = 'paid';

-- one row per group
select status, count(*) as n, sum(total) from tenant_id.table_name group by status;
select status, avg(total) from tenant_id.table_name group by status order by avg(total) desc limit 3;
```

The aggregates are `count(*)`, `count(field)`, `sum`, `min`, `max` and
`avg`, and `AS` names one. As in SQL, `count(*)` counts rows while every
other aggregate skips NULLs, and NULLs form one group. `sum` and `avg` take
numbers and `Decimal`s; `min` and `max` also work on strings and datetimes.

A selected field must be in `GROUP BY`. `ORDER BY` can name a group field,
an aggregate or its alias, and `LIMIT` limits the groups. `WHERE` is applied
before grouping and gets the usual index and primary key pushdown; the
grouping itself happens after the rows are read. `HAVING` is not supported
yet.

### Limit

```sql
Expand Down
132 changes: 132 additions & 0 deletions lib/efsql/aggregate.ex
Original file line number Diff line number Diff line change
@@ -0,0 +1,132 @@
defmodule Efsql.Aggregate do
@moduledoc """
`GROUP BY` and aggregates over pulled rows, for `Efsql.Executor`'s
`{:aggregate, group_by, aggregates}` operator.

Rows are grouped by the values of the `group_by` fields; equal values
group together even when their terms differ (`Decimal` 1.0 and 1.00, or
datetimes stored at different precisions), and NULLs form one group, as
in SQL. Each group becomes one row holding its group fields and one
column per aggregate. With no group fields, every row is one group, so
an empty input still gives one row (`count(*)` of nothing is 0).

SQL semantics: `count(*)` counts rows, every other aggregate skips NULLs,
and `sum`, `min`, `max` and `avg` of no values are NULL. `sum` and `avg`
take numbers and `Decimal`s, and are `Decimal` when any input is.
`min` and `max` order with `Efsql.Types.compare/2`, so they work on
strings and datetimes too. Groups come out ordered by their key.
"""

alias Efsql.Exception.Unsupported
alias Efsql.Types

@type aggregate :: {name :: atom, :count | :sum | :min | :max | :avg, atom | :star}

@spec run([map], [atom], [aggregate]) :: [map]
def run(rows, group_by, aggregates) do
rows
|> Enum.group_by(fn row -> Enum.map(group_by, &key(Map.get(row, &1))) end)
|> ensure_one_group(group_by)
|> Enum.map(fn {_key, rows} ->
first = List.first(rows, %{})
group = Map.new(group_by, &{&1, Map.get(first, &1)})

Enum.reduce(aggregates, group, fn {name, function, arg}, acc ->
Map.put(acc, name, compute(function, arg, rows))
end)
end)
|> Enum.sort(&(compare_groups(&1, &2, group_by) != :gt))
end

# A whole-table aggregate answers even when no row matched.
defp ensure_one_group(groups, []) when groups == %{}, do: [{[], []}]
defp ensure_one_group(groups, _group_by), do: Enum.to_list(groups)

# A term that is equal for values Types.compare/2 calls equal.
defp key(%Decimal{} = d), do: Decimal.normalize(d)
defp key(%NaiveDateTime{microsecond: {us, _}} = t), do: %{t | microsecond: {us, 6}}
defp key(%DateTime{microsecond: {us, _}} = t), do: %{t | microsecond: {us, 6}}
defp key(%Time{microsecond: {us, _}} = t), do: %{t | microsecond: {us, 6}}
defp key(value), do: value

# -- aggregates --

defp compute(:count, :star, rows), do: length(rows)
defp compute(:count, field, rows), do: length(values(rows, field))

defp compute(:min, field, rows), do: extreme(values(rows, field), :lt)
defp compute(:max, field, rows), do: extreme(values(rows, field), :gt)

defp compute(:sum, field, rows) do
case values(rows, field) do
[] -> nil
values -> sum(values, field)
end
end

defp compute(:avg, field, rows) do
case values(rows, field) do
[] ->
nil

values ->
case sum(values, field) do
%Decimal{} = total -> Decimal.div(total, length(values))
total -> total / length(values)
end
end
end

defp values(rows, field) do
rows |> Enum.map(&Map.get(&1, field)) |> Enum.reject(&is_nil/1)
end

defp extreme([], _keep), do: nil

defp extreme([first | rest], keep) do
Enum.reduce(rest, first, fn value, best ->
if Types.compare(value, best) == keep, do: value, else: best
end)
end

defp sum(values, field) do
Enum.each(values, fn value ->
unless is_number(value) or is_struct(value, Decimal) do
raise Unsupported,
"sum and avg need numbers, but #{field} holds #{inspect(value, limit: 5)}"
end
end)

if Enum.any?(values, &is_struct(&1, Decimal)),
do: values |> Enum.map(&decimal/1) |> Enum.reduce(&Decimal.add/2),
else: Enum.sum(values)
end

defp decimal(%Decimal{} = d), do: d
defp decimal(n) when is_integer(n), do: Decimal.new(n)
defp decimal(n) when is_float(n), do: Decimal.from_float(n)

# -- group order --

# By key, NULLs last, like an ascending ORDER BY.
defp compare_groups(_a, _b, []), do: :eq

defp compare_groups(a, b, [field | rest]) do
case {Map.get(a, field), Map.get(b, field)} do
{nil, nil} ->
compare_groups(a, b, rest)

{nil, _} ->
:gt

{_, nil} ->
:lt

{va, vb} ->
case Types.compare(va, vb) do
:eq -> compare_groups(a, b, rest)
cmp -> cmp
end
end
end
end
19 changes: 11 additions & 8 deletions lib/efsql/complete.ex
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ defmodule Efsql.Complete do
"""

@statement_start ~w(select)
@post_expr ~w(and order limit)
@post_expr ~w(and group order limit)
@operators ~w(= > >= < <= like in between not is)

def complete(input, context) do
Expand Down Expand Up @@ -53,14 +53,14 @@ defmodule Efsql.Complete do
last in ["where", "and", "not"] ->
fields(context, tokens)

last == "order" ->
last in ["order", "group"] ->
~w(by)

last == "by" ->
fields(context, tokens)

last in ["asc", "desc"] ->
@post_expr -- ["order"]
~w(limit)

last in @operators ->
[]
Expand All @@ -69,8 +69,9 @@ defmodule Efsql.Complete do
true ->
case section(tokens) do
:select -> ~w(from)
:from -> ~w(where order limit)
:from -> ~w(where group order limit)
:where -> @operators
:group_by -> ~w(order limit)
:order_by -> ~w(asc desc limit)
_ -> []
end
Expand All @@ -95,11 +96,13 @@ defmodule Efsql.Complete do
defp section(tokens) do
tokens
|> Enum.reverse()
|> Enum.chunk_every(2, 1)
|> Enum.find_value(:start, fn
"by" -> :order_by
"where" -> :where
"from" -> :from
"select" -> :select
["by", "group"] -> :group_by
["by" | _] -> :order_by
["where" | _] -> :where
["from" | _] -> :from
["select" | _] -> :select
_ -> nil
end)
end
Expand Down
4 changes: 4 additions & 0 deletions lib/efsql/executor.ex
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,10 @@ defmodule Efsql.Executor do
Enum.filter(rows, fn row -> Enum.all?(predicates, &eval(&1, row)) end)
end

defp apply_op({:aggregate, group_by, aggregates}, rows) do
Efsql.Aggregate.run(rows, group_by, aggregates)
end

defp apply_op({:sort, sort}, rows) do
Enum.sort(rows, fn a, b -> compare(a, b, sort) != :gt end)
end
Expand Down
12 changes: 11 additions & 1 deletion lib/efsql/logical.ex
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,12 @@ defmodule Efsql.Logical do

The primary key is the pseudo-field `:_`.

A grouped query (`GROUP BY`, or any aggregate) has `group_by` set to its
key fields (`[]` for one group over every row) and `aggregates` to
`{output_name, function, field | :star}`, `function` one of `:count`,
`:sum`, `:min`, `:max`, `:avg`. Its `projection` and `order` then name
output columns: group fields and aggregate names.

`Efsql.Rewrite` normalizes a logical query, `Efsql.Planner` turns it into
an `Efsql.Physical.Plan`.
"""
Expand All @@ -29,7 +35,11 @@ defmodule Efsql.Logical do
predicates: [],
# [{:asc | :desc, field}]
order: [],
limit: nil
limit: nil,
# nil, or the GROUP BY fields of a grouped query
group_by: nil,
# [{output_name :: atom, function :: atom, field :: atom | :star}]
aggregates: []
end

def predicate_field({:cmp, _op, field, _value}), do: field
Expand Down
113 changes: 110 additions & 3 deletions lib/efsql/parser.ex
Original file line number Diff line number Diff line change
Expand Up @@ -24,19 +24,126 @@ defmodule Efsql.Parser do
def to_logical(%AST.Select{} = select) do
{prefix, source} = split_from(select.from)

%Logical.Select{
projection: projection(select.fields),
logical = %Logical.Select{
source: source,
prefix: prefix,
predicates: predicates(select.where),
order: Enum.map(select.order_by, fn {name, dir} -> {dir, field(name)} end),
limit: select.limit
}

if grouped?(select),
do: grouped(logical, select),
else: %Logical.Select{
logical
| projection: projection(select.fields),
order: Enum.map(select.order_by, fn {name, dir} -> {dir, field(name)} end)
}
end

defp projection(:star), do: :star
defp projection(fields), do: Enum.map(fields, &field/1)

# -- GROUP BY and aggregates --

@aggregates ~w[count sum min max avg]

defp grouped?(%AST.Select{fields: fields, group_by: group_by, order_by: order_by}) do
group_by != [] or
(is_list(fields) and Enum.any?(fields, &aggregate?/1)) or
Enum.any?(order_by, fn {key, _dir} -> aggregate?(key) end)
end

defp aggregate?(item), do: match?({:aggregate, _, _, _}, item)

# Every selected name must be a group field; ORDER BY may also name an
# aggregate or its alias. An aggregate only in ORDER BY is computed, then
# projected away.
defp grouped(_logical, %AST.Select{fields: :star}) do
raise Unsupported, "SELECT * can't be used with GROUP BY or aggregates; name the fields"
end

defp grouped(logical, select) do
group_by = Enum.map(select.group_by, &group_field/1)

{columns, aggregates} =
Enum.map_reduce(select.fields, [], fn
{:aggregate, _, _, _} = call, aggregates ->
{name, aggregate} = aggregate(call)
{name, aggregates ++ [aggregate]}

name, aggregates ->
field = field(name)

unless field in group_by do
raise Unsupported,
"#{name} must be in GROUP BY or used in an aggregate (count, sum, ...)"
end

{field, aggregates}
end)

{order, aggregates} =
Enum.map_reduce(select.order_by, aggregates, fn
{{:aggregate, _, _, _} = call, dir}, aggregates ->
{name, {_, function, arg} = aggregate} = aggregate(call)

case Enum.find(aggregates, &match?({_, ^function, ^arg}, &1)) do
{existing, _, _} -> {{dir, existing}, aggregates}
nil -> {{dir, name}, aggregates ++ [aggregate]}
end

{name, dir}, aggregates ->
field = field(name)
named = Enum.map(aggregates, &elem(&1, 0))

unless field in group_by or field in named do
raise Unsupported,
"ORDER BY #{name} must name a GROUP BY field, an aggregate or its alias"
end

{{dir, field}, aggregates}
end)

%Logical.Select{
logical
| projection: columns,
order: order,
group_by: group_by,
aggregates: aggregates
}
end

defp group_field("_"),
do: raise(Unsupported, "GROUP BY '_' is not supported; use the primary key field name")

defp group_field(name), do: field(name)

defp aggregate({:aggregate, function, arg, alias}) do
unless function in @aggregates do
raise Unsupported,
"#{function}() is not supported; the aggregates are #{Enum.join(@aggregates, ", ")}"
end

arg =
case arg do
:star when function == "count" ->
:star

:star ->
raise Unsupported, "only count takes *; #{function} needs a field"

"_" ->
raise Unsupported,
"#{function}('_') is not supported; use the primary key field name"

name ->
field(name)
end

name = field(alias || "#{function}(#{if arg == :star, do: "*", else: arg})")
{name, {name, String.to_atom(function), arg}}
end

defp split_from([table]), do: {nil, table}
defp split_from([tenant, table]), do: {tenant, table}
defp split_from([storage, tenant, table]), do: {{storage, tenant}, table}
Expand Down
Loading
Loading