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 @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