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:

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

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.