Packages
ecto_ch
0.3.1
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
Current section
Files
Jump to
Current section
Files
lib/ecto/adapters/clickhouse/structure.ex
defmodule Ecto.Adapters.ClickHouse.Structure do
@moduledoc false
alias Ch.Query
alias Ch.Connection, as: Conn
@conn Ecto.Adapters.ClickHouse.Connection
def structure_load(default, config) do
path = config[:dump_path] || Path.join(default, "structure.sql")
with {:ok, conn} <- Conn.connect(config),
{:ok, queries} <- File.read(path) do
multiquery_result =
queries
|> String.split(";", trim: true)
|> Enum.map(&String.trim/1)
|> Enum.reject(&(&1 == ""))
|> Enum.reduce_while({:ok, _prev_result = nil, conn}, fn
query, {:ok, _prev_result, conn} -> {:cont, exec(conn, query)}
_query, {:error, _reason} = error -> {:halt, error}
end)
case multiquery_result do
{:ok, _last_result, _conn} -> {:ok, path}
{:error, reason} -> {:error, Exception.message(reason)}
end
end
end
# TODO include views
def structure_dump(default, config) do
path = config[:dump_path] || Path.join(default, "structure.sql")
migration_source = config[:migration_source] || "schema_migrations"
database = config[:database] || "default"
with {:ok, conn} <- Conn.connect(config),
{:ok, tables, conn} <- show("TABLES", conn),
{:ok, dicts, conn} <- show("DICTIONARIES", conn),
tables = tables -- [migration_source],
{:ok, tables, conn} <- show_create("TABLE", conn, [migration_source | tables]),
{:ok, dicts, conn} <- show_create("DICTIONARY", conn, dicts),
{:ok, versions, _conn} <- dump_versions(conn, database, migration_source) do
File.mkdir_p!(Path.dirname(path))
File.write!(path, [tables, dicts, versions])
{:ok, path}
end
end
defp show(what, conn) do
with {:ok, %{rows: rows}, conn} <- exec(conn, "SHOW #{what}") do
objects = Enum.map(rows, fn [object] -> object end)
{:ok, objects, conn}
end
end
defp show_create(what, conn, objects) do
show = fn object -> "SHOW CREATE #{what} #{@conn.quote_name(object)}" end
result =
Enum.reduce_while(objects, {[], conn}, fn object, {schemas, conn} ->
case exec(conn, show.(object)) do
{:ok, %{rows: [[schema]]}, conn} -> {:cont, {[schema, ";\n\n" | schemas], conn}}
{:error, _reason} = error -> {:halt, error}
end
end)
case result do
{:error, _reason} = error -> error
{schemas, conn} when is_list(schemas) -> {:ok, schemas, conn}
end
end
defp dump_versions(conn, database, table) do
table = @conn.quote_table(database, table)
stmt = "SELECT * FROM #{table} FORMAT Values"
with {:ok, %{rows: rows}, conn} <- exec(conn, stmt) do
rows = rows |> IO.iodata_to_binary() |> String.replace("),(", "),\n(")
versions = ["INSERT INTO ", table, " (version, inserted_at) VALUES\n", rows, ";\n"]
{:ok, versions, conn}
end
end
def exec(conn, sql, params \\ [], opts \\ []) do
query = Query.build(sql)
params = DBConnection.Query.encode(query, params, [])
case Conn.handle_execute(query, params, opts, conn) do
{:ok, query, result, conn} -> {:ok, DBConnection.Query.decode(query, result, []), conn}
{:disconnect, reason, _conn} -> {:error, reason}
{:error, reason, _conn} -> {:error, reason}
end
end
end