Current section
Files
Jump to
Current section
Files
lib/sink/sql.ex
defmodule Exd.Sink.SQL do
@moduledoc """
A sink that inserts output into a SQL-compatible datastore
## Options
* `:table_name` - Name of table in which to store sink results
* `:key_column_name` - Name of column in which to store document key
* `:data_column_name` - Name of column in which to store document data
## Usage
Query.new()
|> Query.from("jobs", {...})
|> Query.where("jobs.salary", :>, 10000)
|> Query.into(
Exd.Sink.SQL,
hostname: "localhost",
database: "my-database",
username: "postgres",
password: "postgres",
datakey: "jobs.title"
)
"""
use Exd.Sink.Adapter, name: :sql
@default_table_name "documents"
@default_key_column_name "key"
@default_data_column_name "data"
defstruct [
:table_name,
:key_column_name,
:data_column_name,
:connection,
:datakey
]
@impl true
def init(opts) do
{:ok, connection} = init_connection(opts)
table_name = Keyword.get(opts, :table_name, @default_table_name)
key_column_name = Keyword.get(opts, :key_column_name, @default_key_column_name)
data_column_name = Keyword.get(opts, :data_column_name, @default_data_column_name)
datakey = Keyword.fetch!(opts, :datakey)
state =
%__MODULE__{
table_name: table_name,
key_column_name: key_column_name,
data_column_name: data_column_name,
connection: connection,
datakey: datakey
}
:ok = prepare_table(state)
{:ok, state}
end
defp init_connection(opts) do
hostname = Keyword.fetch!(opts, :hostname)
database = Keyword.fetch!(opts, :database)
username = Keyword.fetch!(opts, :username)
password = Keyword.fetch!(opts, :password)
Postgrex.start_link(
hostname: hostname,
database: database,
username: username,
password: password
)
end
@impl true
def handle_into(documents, state) do
{:ok, _results} = insert_values(documents, state)
{:ok, state}
end
defp prepare_table(state) do
create_table_query = "
CREATE TABLE IF NOT EXISTS #{state.table_name} (
ID serial NOT NULL PRIMARY KEY,
#{state.key_column_name} varchar(100) NOT NULL,
#{state.data_column_name} json NOT NULL
);
"
create_unique_index_query = "
CREATE UNIQUE INDEX IF NOT EXISTS #{state.key_column_name}_idx ON #{state.table_name} (#{state.key_column_name});
"
with {:ok, _} <- Postgrex.query(state.connection, create_table_query, []),
{:ok, _} <- Postgrex.query(state.connection, create_unique_index_query, []) do
:ok
end
end
defp insert_values(documents, state) do
values =
for document <- documents do
document_key = Map.fetch!(document, state.datakey)
"('#{document_key}', '#{Poison.encode!(document)}')"
end
|> Enum.join(", ")
query = "
INSERT INTO #{state.table_name} (#{state.key_column_name}, #{state.data_column_name})
VALUES #{values}
ON CONFLICT (#{state.key_column_name})
DO UPDATE SET #{state.data_column_name} = EXCLUDED.#{state.data_column_name};
"
Postgrex.query(state.connection, query, [])
end
end