Packages
pillar
0.16.1
0.40.0
0.39.0
0.38.0
0.37.0
0.36.0
0.35.0
0.34.1
0.34.0
0.33.1
0.33.0
0.32.0
0.31.0
0.30.0
0.29.1
0.29.0
0.28.0
0.27.0
0.26.1
0.26.0
0.25.1
0.25.0
0.24.0
0.23.3
0.23.2
0.23.1
0.23.0
0.22.0
0.21.0
0.20.0
0.19.0
0.18.2
0.18.1
0.18.0
0.17.3
0.17.2
0.17.1
0.17.0
0.16.2
0.16.1
0.16.0
0.15.0
0.14.0
0.13.1
0.13.0
0.12.0
0.11.0
0.10.0
0.9.1
0.9.0
0.8.1
0.8.0
0.7.0
0.6.0
0.5.1
0.5.0
0.4.0
0.3.2
0.3.1
0.3.0
0.2.1
0.2.0
0.1.0
Elixir client for ClickHouse, a fast open-source Online Analytical Processing (OLAP) database management system.
Current section
Files
Jump to
Current section
Files
lib/pillar/bulk_insert_buffer.ex
defmodule Pillar.BulkInsertBuffer do
@moduledoc """
This module provides functionality for bulk inserts and buffering records
```elixir
defmodule BulkToLogs do
use Pillar.BulkInsertBuffer,
pool: ClickhouseMaster,
table_name: "logs",
interval_between_inserts_in_seconds: 5
end
```
```elixir
:ok = BulkToLogs.insert(%{value: "online", count: 133, datetime: DateTime.utc_now()})
:ok = BulkToLogs.insert(%{value: "online", count: 134, datetime: DateTime.utc_now()})
:ok = BulkToLogs.insert(%{value: "online", count: 132, datetime: DateTime.utc_now()})
....
# all this records will be inserted with 5 second interval
```
"""
alias Pillar.TypeConvert.ToClickhouseJson
def generate_insert_query(table_name, values) do
new_values = Enum.map(values, &convert_values_to_clickhouse/1)
sql_strings = [
"INSERT INTO",
table_name,
"FORMAT JSONEachRow",
Enum.join(Enum.map(new_values, &Jason.encode!/1), " ")
]
Enum.join(sql_strings, "\n")
end
defp convert_values_to_clickhouse(map) do
map
|> Enum.reject(fn {_key, value} -> is_nil(value) end)
|> Enum.map(fn {key, value} ->
{key, ToClickhouseJson.convert(value)}
end)
|> Map.new()
end
defmacro __using__(
pool: pool_module,
table_name: table_name,
interval_between_inserts_in_seconds: seconds
) do
quote do
use GenServer
import Supervisor.Spec
alias Pillar.BulkInsertBuffer
def start_link(_any \\ nil) do
name = unquote(__MODULE__)
pool = unquote(pool_module)
table_name = unquote(table_name)
records = []
GenServer.start_link(__MODULE__, {pool, table_name, records}, name: name)
end
def init(state) do
schedule_work()
{:ok, state}
end
def insert(data) when is_map(data) do
GenServer.call(unquote(__MODULE__), {:insert, data})
end
def force_bulk_insert do
GenServer.call(unquote(__MODULE__), :do_insert)
end
def records_for_bulk_insert() do
GenServer.call(unquote(__MODULE__), :records_for_bulk_insert)
end
def handle_call(:do_insert, _from, state) do
new_state = do_bulk_insert(state)
{:reply, :ok, new_state}
end
def handle_call({:insert, data}, _from, {pool, table_name, records} = state) do
{:reply, :ok, {pool, table_name, records ++ List.wrap(data)}}
end
def handle_call(
:records_for_bulk_insert,
_from,
{_pool, _table_name, records} = state
) do
{:reply, records, state}
end
def handle_info(:cron_like_records, state) do
new_state = do_bulk_insert(state)
schedule_work()
{:noreply, new_state}
end
defp schedule_work do
seconds = unquote(seconds)
Process.send_after(self(), :cron_like_records, :timer.seconds(seconds))
end
defp do_bulk_insert({_pool, _table_name, []} = state) do
state
end
defp do_bulk_insert({pool, table_name, records} = state) do
sql = BulkInsertBuffer.generate_insert_query(table_name, records)
{:ok, _} = pool.query(sql, %{})
{
pool,
table_name,
[]
}
end
end
end
end