Packages

Bulk insert and file-driven sync for Ecto via unnest(...) — pass data as columns (%{col => list}) rather than rows. Constant SQL text independent of row count, friendly to PgBouncer (transaction mode) and the prepared-statement cache.

Current section

Files

Jump to
ecto_unnest README.md
Raw

README.md

# 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:
```sql
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
```sql
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:
```elixir
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](#syncing-a-table-from-a-file-sync_all4)).
`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](#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:
```elixir
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`:
```elixir
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:
```csv
id,val
1,x
2,y
3,z
```
```elixir
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.
```elixir
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:
```sql
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.
```elixir
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:
```elixir
EctoUnnest.sync_all(Repo, Event, %{user_id: [1, 2, 3], type: ["a", "b", "c"]},
key: [:user_id])
```
```sql
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.
```elixir
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:
```elixir
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:
```elixir
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`](https://hex.pm/packages/uuuidv7)
(`{:uuuidv7, "~> 0.3.0"}`):
```elixir
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.