Packages
ecto_ch
0.4.0
0.10.0
0.9.4
0.9.3
0.9.2
0.9.1
0.9.0
0.8.9
0.8.8
0.8.7
0.8.6
0.8.5
retired
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
retired
0.6.1
retired
0.6.0
retired
0.5.1
0.5.0
retired
0.4.1
0.4.0
retired
0.3.10
0.3.9
0.3.8
0.3.7
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.3.1
0.3.0
0.2.2
0.2.1
0.2.0
0.1.11
0.1.10
0.1.9
0.1.8
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
ClickHouse adapter for Ecto
Retired package: Deprecated - Partially incompatible with Ecto v3.13. Use v0.4.1 or v0.6.3 for Ecto v3.12, or v0.7.0 for Ecto v3.13.
Current section
Files
Jump to
Current section
Files
lib/ecto/adapters/clickhouse.ex
defmodule Ecto.Adapters.ClickHouse do
@moduledoc """
Adapter module for ClickHouse.
It uses `Ch` for communicating to the database.
## Options
All options can be given via the repository
configuration:
config :your_app, YourApp.Repo,
...
* `:hostname` - Server hostname (default: `"localhost"`)
* `:username` - Username
* `:password` - User password
* `:port` - HTTP Server port (default: `8123`)
* `:scheme` - HTTP scheme (default: `"http"`)
* `:database` - the database to connect to (default: `"default"`)
* `:settings` - Keyword list of connection settings
* `:transport_opts` - Options to be given to the transport being used. See `Mint.HTTP1.connect/4` for more info
"""
@behaviour Ecto.Adapter
@behaviour Ecto.Adapter.Migration
@behaviour Ecto.Adapter.Queryable
@behaviour Ecto.Adapter.Schema
@behaviour Ecto.Adapter.Storage
@behaviour Ecto.Adapter.Structure
@conn __MODULE__.Connection
@driver :ch
@impl Ecto.Adapter
defmacro __before_compile__(_env) do
quote do
@doc """
A convenience function for SQL-based repositories that executes the given query.
See `Ecto.Adapters.SQL.query/4` for more information.
"""
def query(sql, params \\ [], opts \\ []) do
Ecto.Adapters.SQL.query(get_dynamic_repo(), sql, params, opts)
end
@doc """
A convenience function for SQL-based repositories that executes the given query.
See `Ecto.Adapters.SQL.query!/4` for more information.
"""
def query!(sql, params \\ [], opts \\ []) do
Ecto.Adapters.SQL.query!(get_dynamic_repo(), sql, params, opts)
end
@doc """
A convenience function for SQL-based repositories that translates the given query to SQL.
See `Ecto.Adapters.SQL.to_sql/3` for more information.
"""
def to_sql(operation, queryable) do
Ecto.Adapters.ClickHouse.to_sql(operation, queryable)
end
@doc """
Similar to `to_sql/2` but inlines the parameters into the SQL query.
See `Ecto.Adapters.ClickHouse.to_inline_sql/2` for more information.
"""
@spec to_inline_sql(:all | :delete_all | :update_all, Ecto.Queryable.t()) :: String.t()
def to_inline_sql(operation, queryable) do
Ecto.Adapters.ClickHouse.to_inline_sql(operation, queryable)
end
@doc """
A convenience function for SQL-based repositories that forces all connections in the
pool to disconnect within the given interval.
See `Ecto.Adapters.SQL.disconnect_all/3` for more information.
"""
def disconnect_all(interval, opts \\ []) do
Ecto.Adapters.SQL.disconnect_all(get_dynamic_repo(), interval, opts)
end
@doc """
Similar to `insert_all/2` but with the following differences:
- accepts rows as streams or lists
- sends rows as a chunked request
- doesn't autogenerate ids or does any other preprocessing
Example:
Repo.query!("create table ecto_ch_demo(a UInt64, b String) engine Null")
defmodule Demo do
use Ecto.Schema
@primary_key false
schema "ecto_ch_demo" do
field :a, Ch, type: "UInt64"
field :b, :string
end
end
rows = Stream.map(1..100_000, fn i -> %{a: i, b: to_string(i)} end)
{100_000, nil} = Repo.insert_stream(Demo, rows)
# schemaless
{100_000, nil} = Repo.insert_stream("ecto_ch_demo", rows, types: [a: Ch.Types.u64(), b: :string])
"""
def insert_stream(source_or_schema, rows, opts \\ []) do
repo = get_dynamic_repo()
# TODO need it?
# opts = Ecto.Repo.Supervisor.tuplet(repo, prepare_opts(:insert_all, opts))
Ecto.Adapters.ClickHouse.Schema.insert_stream(repo, source_or_schema, rows, opts)
end
@doc """
Similar to `Ecto.Repo.update_all/3` but uses [`ALTER TABLE ... UPDATE`](https://clickhouse.com/docs/en/sql-reference/statements/alter/update) instead.
For more information and performance implications please see:
- https://clickhouse.com/blog/handling-updates-and-deletes-in-clickhouse
- https://clickhouse.com/docs/en/guides/developer/mutations
"""
def alter_update_all(queryable, updates, opts \\ []) do
repo = get_dynamic_repo()
Ecto.Adapters.ClickHouse.Queryable.alter_update_all(
repo,
queryable,
updates,
Ecto.Repo.Supervisor.tuplet(repo, prepare_opts(:update_all, opts))
)
end
end
end
@impl Ecto.Adapter
def ensure_all_started(config, type) do
Ecto.Adapters.SQL.ensure_all_started(@driver, config, type)
end
@impl Ecto.Adapter
def init(config) do
Ecto.Adapters.SQL.init(@conn, @driver, config)
end
@impl Ecto.Adapter
def checkout(meta, opts, fun) do
Ecto.Adapters.SQL.checkout(meta, opts, fun)
end
@impl Ecto.Adapter
def checked_out?(meta) do
Ecto.Adapters.SQL.checked_out?(meta)
end
@impl Ecto.Adapter
def dumpers(:uuid, Ecto.UUID), do: [&__MODULE__.hex_uuid/1]
def dumpers(:uuid, type), do: [type, &__MODULE__.hex_uuid/1]
def dumpers(:binary_id, type), do: [type, &__MODULE__.hex_uuid/1]
def dumpers(_primitive, {:parameterized, {Ch, params}}), do: [dumper(params)]
def dumpers({:parameterized, {Ch, params}}, type), do: [type, dumper(params)]
def dumpers(_primitive, type), do: [type]
defp dumper(:uuid), do: &__MODULE__.hex_uuid/1
defp dumper(:date32), do: :date
defp dumper(:datetime), do: :naive_datetime
defp dumper({:datetime, "UTC"}), do: :utc_datetime
defp dumper({:datetime64, _precision}), do: :naive_datetime_usec
defp dumper({:datetime64, _precision, "UTC"}), do: :utc_datetime_usec
defp dumper({:nullable, type}), do: dumper(type)
defp dumper({:low_cardinality, type}), do: dumper(type)
defp dumper({:decimal = d, _precision, _scale}), do: d
for size <- [32, 64, 128, 256] do
defp dumper({unquote(:"decimal#{size}"), _scale}), do: :decimal
end
for size <- [8, 16, 32, 64, 128, 256] do
defp dumper(unquote(:"i#{size}")), do: :integer
defp dumper(unquote(:"u#{size}")), do: :integer
end
for size <- [32, 64] do
defp dumper(unquote(:"f#{size}")), do: :float
end
defp dumper({:simple_aggregate_function, _name, type}), do: dumper(type)
defp dumper({:array = array, type}), do: {array, dumper(type)}
defp dumper(_type), do: &__MODULE__.ok_identity/1
@impl Ecto.Adapter
def loaders(:uuid, Ecto.UUID = uuid), do: [uuid]
def loaders(:uuid, type), do: [Ecto.UUID, type]
def loaders(:binary_id, type), do: [Ecto.UUID, type]
def loaders(_primitive, {:parameterized, {Ch, params}}), do: [loader(params)]
def loaders({:parameterized, {Ch, params}}, type), do: [loader(params), type]
def loaders(_primitive, type), do: [type]
defp loader(:uuid), do: Ecto.UUID
defp loader({:nullable, type}), do: loader(type)
defp loader({:low_cardinality, type}), do: loader(type)
defp loader({:simple_aggregate_function, _name, type}), do: loader(type)
defp loader({:array = array, type}), do: {array, loader(type)}
defp loader(_type), do: &__MODULE__.ok_identity/1
@doc false
def ok_identity(value), do: {:ok, value}
@doc false
def hex_uuid(nil), do: {:ok, nil}
def hex_uuid(value), do: Ecto.UUID.cast(value)
@impl Ecto.Adapter.Migration
def supports_ddl_transaction?, do: false
@impl Ecto.Adapter.Migration
def lock_for_migrations(_meta, _options, f), do: f.()
@impl Ecto.Adapter.Migration
def execute_ddl(meta, definition, opts) do
Ecto.Adapters.SQL.execute_ddl(meta, @conn, definition, opts)
end
@impl Ecto.Adapter.Storage
defdelegate storage_up(opts), to: Ecto.Adapters.ClickHouse.Storage
@impl Ecto.Adapter.Storage
defdelegate storage_down(opts), to: Ecto.Adapters.ClickHouse.Storage
@impl Ecto.Adapter.Storage
defdelegate storage_status(opts), to: Ecto.Adapters.ClickHouse.Storage
@impl Ecto.Adapter.Structure
defdelegate structure_dump(default, config), to: Ecto.Adapters.ClickHouse.Structure
@impl Ecto.Adapter.Structure
defdelegate structure_load(default, config), to: Ecto.Adapters.ClickHouse.Structure
@impl Ecto.Adapter.Structure
def dump_cmd(_args, _opts, _config) do
raise "not implemented"
end
@impl Ecto.Adapter.Schema
def autogenerate(:id), do: nil
def autogenerate(:embed_id), do: Ecto.UUID.generate()
def autogenerate(:binary_id), do: Ecto.UUID.bingenerate()
@impl Ecto.Adapter.Schema
def insert_all(
adapter_meta,
schema_meta,
header,
rows,
on_conflict,
returning,
placeholders,
opts
) do
Ecto.Adapters.ClickHouse.Schema.insert_all(
adapter_meta,
schema_meta,
header,
rows,
on_conflict,
returning,
placeholders,
opts
)
end
@impl Ecto.Adapter.Schema
def insert(adapter_meta, schema_meta, params, _, _, opts) do
Ecto.Adapters.ClickHouse.Schema.insert(adapter_meta, schema_meta, params, opts)
end
@dialyzer {:no_return, update: 6}
@impl Ecto.Adapter.Schema
def update(adapter_meta, %{source: source, prefix: prefix}, fields, params, returning, opts) do
{fields, field_values} = :lists.unzip(fields)
filter_values = Keyword.values(params)
sql = @conn.update(prefix, source, fields, params, returning)
Ecto.Adapters.SQL.struct(
adapter_meta,
@conn,
sql,
:update,
source,
params,
field_values ++ filter_values,
:raise,
returning,
opts
)
end
# TODO
# https://github.com/elixir-ecto/ecto/blob/master/CHANGELOG.md#v3110-2023-11-14
# https://github.com/elixir-ecto/ecto/pull/4277
# @impl Ecto.Adapter.Schema
def delete(adapter_meta, schema_meta, params, opts) do
Ecto.Adapters.ClickHouse.Schema.delete(adapter_meta, schema_meta, params, opts)
end
@impl Ecto.Adapter.Schema
def delete(adapter_meta, schema_meta, params, _returning, opts) do
Ecto.Adapters.ClickHouse.Schema.delete(adapter_meta, schema_meta, params, opts)
end
@impl Ecto.Adapter.Queryable
def stream(adapter_meta, query_meta, query, params, opts) do
Ecto.Adapters.SQL.stream(adapter_meta, query_meta, query, params, opts)
end
@impl Ecto.Adapter.Queryable
def prepare(operation, query), do: {:nocache, {operation, query}}
@impl Ecto.Adapter.Queryable
def execute(adapter_meta, query_meta, {:nocache, {operation, query}}, params, opts) do
sql = prepare_sql(operation, query, params)
opts =
case operation do
:all -> [{:command, :select} | opts]
:delete_all -> [{:command, :delete} | opts]
:alter_update_all -> [{:command, :alter} | opts]
end
result = Ecto.Adapters.SQL.query!(adapter_meta, sql, params, put_source(opts, query_meta))
case operation do
:all ->
%{num_rows: num_rows, rows: rows} = result
{num_rows, rows}
:delete_all ->
# clickhouse doesn't give us any info on how many rows have been deleted
{0, nil}
:alter_update_all ->
# clickhouse doesn't give us any info on how many rows have been / will be altered
{0, nil}
end
end
@doc false
def to_sql(operation, queryable) do
{query, _cast_params, dump_params} = plan_query(operation, queryable)
sql = Ecto.Adapters.ClickHouse.prepare_sql(operation, query, dump_params)
{IO.iodata_to_binary(sql), dump_params}
end
@doc """
Converts the given query to SQL according to its kind and inlines all parameters. Useful for debugging.
Example:
iex> query = from n in fragment("numbers(10)"), where: n.number == ^-1, select: n.number
iex> Ecto.Adapters.ClickHouse.to_inline_sql(:all, query)
~s{SELECT f0."number" FROM numbers(10) AS f0 WHERE (f0."number" = -1)}
Compare it with the output of `to_sql/2`
iex> query = from n in fragment("numbers(10)"), where: n.number == ^-1, select: n.number
iex> Ecto.Adapters.ClickHouse.to_sql(:all, query)
{~s[SELECT f0."number" FROM numbers(10) AS f0 WHERE (f0."number" = {$0:Int64})], [-1]}
"""
@spec to_inline_sql(:all | :delete_all | :update_all, Ecto.Queryable.t()) :: String.t()
def to_inline_sql(operation, queryable) do
{query, _cast_params, dump_params} = plan_query(operation, queryable)
inline_params = Enum.map(dump_params, &@conn.mark_inline/1)
sql = Ecto.Adapters.ClickHouse.prepare_sql(operation, query, inline_params)
IO.iodata_to_binary(sql)
end
defp plan_query(operation, queryable) do
queryable =
queryable
|> Ecto.Queryable.to_query()
|> Ecto.Query.Planner.ensure_select(operation == :all)
operation =
case operation do
:alter_update_all -> :update_all
:alter_delete_all -> :delete_all
_other -> operation
end
Ecto.Adapter.Queryable.plan_query(operation, Ecto.Adapters.ClickHouse, queryable)
end
@doc false
def prepare_sql(:all, query, params), do: @conn.all(query, params)
def prepare_sql(:update_all, query, params), do: @conn.update_all(query, params)
def prepare_sql(:delete_all, query, params), do: @conn.delete_all(query, params)
def prepare_sql(:alter_update_all, query, params), do: @conn.alter_update_all(query, params)
# TODO
# def prepare_sql(:alter_delete_all, query, params), do: @conn.alter_delete_all(query, params)
defp put_source(opts, %{sources: sources}) when is_binary(elem(elem(sources, 0), 0)) do
{source, _, _} = elem(sources, 0)
[source: source] ++ opts
end
defp put_source(opts, _) do
opts
end
end