EctoUnnest
Bulk insert for Ecto via unnest(...) — constant SQL text, independent of the
row count, friendly to PgBouncer (transaction mode) and the prepared-statement
cache.
Problem
Ecto.Repo.insert_all/3 builds:
INSERT INTO events (a, b) VALUES ($1, $2), ($3, $4), ... -- N*K parameters
Every batch size is a different SQL text → a different prepared statement.
PgBouncer in transaction mode can't cache that sensibly, and Postgres caps
parameters at ~65535.
Solution
INSERT INTO "events" ("type","user_id")
(SELECT f0."type", f0."user_id"
FROM (SELECT * FROM unnest($1::text[], $2::bigint[]) AS u("type","user_id")) AS f0)
Always K parameters (one array per column) → the statement is identical for 1 and 10,000 rows.
The query is assembled from Ecto building blocks (fragment/dynamic) and handed
to Ecto.Repo.insert_all/3, which renders ON CONFLICT/RETURNING/prefix and
loads structs natively. Because the text is constant per shape, Postgres reuses a
single prepared statement (Ecto caches it under ecto_insert_all_<table>).
Usage
Two disjoint maps:
EctoUnnest.insert_all(Repo, Event,
# columns map: each value is a list -> goes into unnest
%{user_id: [1, 2, 3], type: ["click", "view", "click"]},
# :placeholders: constants broadcast onto every row
placeholders: %{inserted_at: ~U[2026-06-17 10:00:00Z]},
returning: true
)
# => {3, [%Event{...}, %Event{...}, %Event{...}]}
EctoUnnest.sync_all/4 goes one step further: it makes the table match the
source, deleting the rows the source no longer has (see
Syncing a table from a file).
EctoUnnest.to_sql/3 returns {sql, params} without executing — for debugging.
It is pure (no database connection) and renders the exact statement
Repo.insert_all/3 would run.
Column names, table names, schema prefixes, conflict targets, and type overrides
are SQL structure. Keep them defined by trusted application code or configuration;
do not build them from user-controlled input. Row values belong in the lists or
:placeholders, where Ecto/Postgrex sends them as parameters.
Options (same as Ecto.Repo.insert_all/3)
| option | meaning |
|---|---|
:placeholders |
%{col => value} of constant columns (default %{}) |
:returning |
true | false | [field] |
:prefix |
schema prefix (overrides @schema_prefix) |
:on_conflict |
:raise | :nothing | :replace_all | {:replace, fields} | {:replace_all_except, fields} | [set: kw, inc: kw] |
:conflict_target |
[col] | {:unsafe_fragment, binary} |
:types |
%{col => pg_type} override for inference (atom or string — see Type overrides) |
:cache_statement |
prepared-statement cache name (default "ecto_unnest_all_#{table}_#{arity}", where arity is the number of unnest columns; pass a binary to override, or nil for Ecto's default) |
:require_all_fields |
true to assert every schema field is in the columns map or :placeholders (defaults to config :ecto_unnest, :require_all_fields, else false) |
Type overrides and :allowed_types
A :types value is rendered into the SQL cast ($n::type / ::type[])
verbatim, so it must never come from user input. As a guard, the gate is
fail-closed:
- Ecto's default PG types (
bigint,float8,boolean,text,bytea,uuid,numeric,date,time,timestamp,timestamptz,jsonb,json) are always accepted. - Anything else — a domain (
kafka_topic_name), an alias (int4), a modifier (numeric(10,2)), a quoted identifier — must be vouched for in config, or the call raises:
config :ecto_unnest, :allowed_types, [:int4, "kafka_topic_name"]
Entries may be atoms or strings. Because binary sources ("table") require
:types, any non-default type they cast to must be listed here.
Reading: unnest as a virtual table
EctoUnnest.table/3 exposes the same unnest(...) source as a composable
%Ecto.Query{} with a named binding (:s by default). Use it like any Ecto
source — where, order_by, select, Repo.all/2:
q = EctoUnnest.table(Event, %{user_id: [1, 2, 3], type: ["a", "b", "c"]})
from([s: s] in q, where: s.user_id > 1, select: {s.user_id, s.type})
|> Repo.all()
# => [{2, "b"}, {3, "c"}]
Bulk UPDATE from an in-memory CSV
Wrap the virtual table in subquery/1 (which carries the parameters) and join it
into an UPDATE. Here we drive the update from a CSV like:
id,val
1,x
2,y
3,z
csv = "id,val\n1,x\n2,y\n3,z\n"
# parse CSV into column arrays: %{id: [1, 2, 3], val: ["x", "y", "z"]}
[_header | rows] = csv |> String.trim() |> String.split("\n")
cols =
rows
|> Enum.map(&String.split(&1, ","))
|> Enum.reduce(%{id: [], val: []}, fn [id, val], acc ->
%{acc | id: [String.to_integer(id) | acc.id], val: [val | acc.val]}
end)
|> Map.update!(:id, &Enum.reverse/1)
|> Map.update!(:val, &Enum.reverse/1)
# build the virtual table and join it into a single UPDATE statement
src =
EctoUnnest.table(Event, %{user_id: cols.id, type: cols.val})
|> then(&from([s: s] in &1, select: %{user_id: s.user_id, type: s.type}))
from(e in Event,
join: s in subquery(src),
on: e.user_id == s.user_id,
update: [set: [type: s.type]]
)
|> Repo.update_all([])
# => {3, nil} — one statement, constant text regardless of CSV size
Syncing a table from a file: sync_all/4
When a file in git — YAML, CSV, JSON — is the source of truth for a table, upsert
alone is not enough: ON CONFLICT only ever sees the rows you sent, so a row you
deleted from the file is a row Postgres never hears about. sync_all/4 does
both halves in one transaction.
settings = YamlElixir.read_from_file!("priv/settings.yml") # [%{"key" => ..., "value" => ...}]
cols = %{
key: Enum.map(settings, & &1["key"]),
value: Enum.map(settings, & &1["value"])
}
EctoUnnest.sync_all(Repo, Setting, cols)
# => %{upserted: 12, deleted: 1, soft_deleted: 0,
# deleted_keys: [%{key: "retired_flag"}], skipped_keys: [], on_delete: :restrict}
Two statements, both with text independent of the row count — the delete reads its
keys back out of the same unnest(...) source:
INSERT INTO "settings" ("value","key")
(SELECT f0."value", f0."key"
FROM (SELECT * FROM unnest($1::text[], $2::text[]) AS u("key", "value")) AS f0)
ON CONFLICT ("key") DO UPDATE SET "value" = EXCLUDED."value"
DELETE FROM "settings" AS s0
WHERE (NOT (exists((SELECT 1 FROM (SELECT * FROM unnest($1::text[]) AS u("key")) AS sf0
WHERE (sf0."key" = s0."key")))))
RETURNING s0."key"
Composite keys need no special casing (key: [:tenant_id, :code] compares both
columns), and :deleted_keys tells you exactly which rows went.
Options
Everything insert_all/4 takes, plus:
| option | meaning |
|---|---|
:key |
column(s) identifying a row (default: the schema's primary key). Must be in the columns map — the file has to say which row it describes |
:on_delete |
what to do with a row the source no longer has (below), default :restrict |
:allow_empty |
true to permit a source with no rows. An empty source means "delete everything", which is nearly always a failed parse, so it raises by default |
:lock |
true to serialize concurrent syncs of this table, or an integer advisory-lock key (default false) |
:conflict_target defaults to :key, and :on_conflict to replacing exactly the
columns the source carries — deliberately not Ecto's :replace_all, which
means every field of the schema and so would overwrite a column the file does not
mention with NULL. Pass a full Ecto.Query as :on_conflict for a conditional
DO UPDATE ... WHERE, so rows whose values did not change do not churn.
:on_delete
Deletion is the only part of a sync that can destroy data the source never
described, because ON DELETE CASCADE reaches rows in other tables — with no
error.
| mode | behaviour |
|---|---|
:restrict (default) |
reads pg_constraint and refuses to run if any foreign key pointing at the table is CASCADE/SET NULL/SET DEFAULT |
:skip_referenced |
deletes only the absent rows nothing references; the rest come back in :skipped_keys |
{:soft, field} |
sets field to the current time instead of deleting. Also revives a tombstoned row that reappears in the source |
{:soft, field, value} |
same, with the value you choose (required for a binary source) |
:cascade |
no catalog check, plain delete — destructive foreign keys fire. Explicit opt-in |
:nothing |
upsert only |
NO ACTION/RESTRICT children need no check: Postgres raises
foreign_key_violation and the transaction rolls the whole sync back, upsert
included. The gap :restrict closes is the silent one.
EctoUnnest.sync_all(Repo, Setting, cols)
** (ArgumentError) sync_all/4 with on_delete: :restrict refuses to delete from "settings":
public.settings_logs(setting_key) ON DELETE CASCADE (settings_logs_setting_key_fkey).
Deleting a row absent from the source would silently change rows in those tables.
Use on_delete: :skip_referenced to leave referenced rows alone, {:soft, field} to
tombstone instead, or :cascade to accept the cascade.
Row identity
A row keeps its surrogate primary key across syncs. Syncing on a natural key is the normal case:
EctoUnnest.sync_all(Repo, Event, %{user_id: [1, 2, 3], type: ["a", "b", "c"]},
key: [:user_id])
INSERT INTO "events" ("type","user_id") (SELECT ... FROM unnest($1::text[], $2::bigint[]) ...)
ON CONFLICT ("user_id") DO UPDATE SET "type" = EXCLUDED."type"
:id is not in the source, so it appears in neither the column list nor the
DO UPDATE SET — existing rows keep the id they had, and removing a row does not
renumber the others.
This is why the default :on_conflict is not Ecto's :replace_all. That
expands to __schema__(:updatable_fields), which includes the primary key, and
for a column the INSERT does not list EXCLUDED holds the column's default —
NULL, or the next sequence value for a serial. So :replace_all on a natural-key
sync renders SET "id" = EXCLUDED."id" and renumbers every row it touches:
id_before | user_id | type id_after | user_id | type
5476 | 7 | orig -> 5477 | 7 | synced
Anything holding a foreign key to that row now points at nothing. So a
caller-supplied :on_conflict is checked: if its replace list names a column the
source does not provide, it raises instead of resetting it.
EctoUnnest.sync_all(Repo, Event, cols, key: [:user_id], on_conflict: :replace_all)
** (ArgumentError) sync_all/4: :on_conflict would replace [:id, :inserted_at, :payload,
:score, :tags], which the source does not provide. `EXCLUDED` holds each column's
default for a column the INSERT does not list, so those would be reset to NULL (or,
for a serial, the next sequence value — renumbering the row and repointing anything
that references it). That includes the primary key. ...
Add the column to the source, exclude it ({:replace_all_except, [:id]}), or leave
:on_conflict alone. The keyword and Ecto.Query forms name their own values rather
than taking them from EXCLUDED, so they are not second-guessed.
Concurrent syncs
A sync is two statements, so two running at once can interleave: both upsert, then
both delete, and the loser's rows are removed by the winner's delete before its own
upsert is visible. lock: true takes a transaction-scoped advisory lock keyed on
the table's oid, so the second sync waits for the first:
EctoUnnest.sync_all(Repo, Setting, cols, lock: true)
# SELECT pg_advisory_xact_lock($1::text::regclass::oid::bigint)
# INSERT INTO "settings" ...
# DELETE FROM "settings" ...
There is nothing to unlock: pg_advisory_xact_lock is released when the
transaction ends, including when the sync raises. It is off by default because it
blocks — turn it on when a sync can fire more than once at a time (a deploy hook
that may overlap, a scheduled job, a web endpoint).
Advisory locks share one key space per database, so true derives the key from the
table oid — unique per table, but in principle able to collide with another
library's key. Pass your own integer to control it, or to take the same lock from
other code:
EctoUnnest.sync_all(Repo, Setting, cols, lock: 8_675_309)
Note the lock lives as long as the outermost transaction, so a sync_all/4
nested inside your own Repo.transaction/1 holds it until yours commits.
For git-tracked reference data {:soft, :deleted_at} is usually the right answer:
it sidesteps foreign keys entirely, and the tombstone value follows the field's
Ecto type (a :utc_datetime column gets a second-truncated DateTime, not
microseconds it cannot store).
Only :restrict and :skip_referenced read the catalog — one small query each.
UUIDv7 primary keys
insert_all (and therefore EctoUnnest) does not autogenerate primary keys, so
you supply the id list yourself. EctoUnnest is type-agnostic: a UUIDv7 field is a
plain Ecto.Type whose type/0 is :uuid, so it is inferred as ::uuid[] and
each id is dumped through that type — no special integration needed.
Pair it with a bulk id generator such as uuuidv7
({:uuuidv7, "~> 0.3.0"}):
names = ["a", "b", "c"]
EctoUnnest.insert_all(Repo, Event,
%{id: UUIDv7.generate_many(length(names)), name: names}
)
UUIDv7.generate_many/1 produces monotonic ids in one shot, which keeps the insert
a single constant-text statement.
Arrays
- A constant array value → insert it via
:placeholders(goes in as a scalar$n::int[], no unnest). Works. - A per-row array column → unsupported (unnest flattens multi-dimensional arrays); the library raises a clear error.
Status
Unit tests (to_sql, table, and sync_all/4's validation) need no database.
Integration tests (tagged @tag :integration) require Postgres; run them with
INTEGRATION=1 mix test. sync_all/4's foreign-key modes are covered there
against real ON DELETE actions.