DML (Data Manipulation Language)
Gnok supports full ACID DML operations on Iceberg tables (format version 2+). Each statement executes as an atomic Iceberg commit with snapshot isolation.
DML operations (UPDATE, DELETE, MERGE) require Iceberg format version 2 or higher — they rely on row-level delete files (position deletes). INSERT is supported on all format versions.
INSERT
Service resource limits are managed by Gnok. Reference defaults in this page describe engine behavior, not account entitlements or settings to export on your client.
INSERT VALUES
-- Single row
INSERT INTO orders (order_id, customer_id, amount, order_date)
VALUES (1, 100, 99.99, '2025-01-15');
-- Multiple rows
INSERT INTO orders (order_id, customer_id, amount, order_date)
VALUES
(1, 100, 99.99, '2025-01-15'),
(2, 101, 49.50, '2025-01-16'),
(3, 102, 199.00, '2025-01-17');
For partitioned tables, Gnok groups rows by partition key before writing, producing optimally-sized data files per partition.
INSERT SELECT
-- Insert from query
INSERT INTO orders_archive
SELECT * FROM orders WHERE order_date < '2024-01-01';
-- With column list and transformations
INSERT INTO summary (region, total_amount, order_count)
SELECT region, SUM(amount), COUNT(*)
FROM orders
GROUP BY region;
Gnok validates schema compatibility between the SELECT output and the target columns, applying automatic type casts where safe (e.g., INT to BIGINT). For partitioned tables, data is shuffled to partition-owning workers for optimal file layout.
INSERT OVERWRITE
Replace data with the new data:
INSERT OVERWRITE orders_daily
SELECT * FROM staging_orders WHERE order_date = '2025-01-15';
On an unpartitioned table, INSERT OVERWRITE is a full-table replace: every prior data manifest is dropped from the new snapshot, and only the rows produced by the SELECT (or the literal VALUES list) survive. The replace runs through the same distributed write path as INSERT … SELECT, so it parallelises across workers and writes optimally-sized Parquet files.
On a partitioned table, only the partitions present in the new data are overwritten — other partitions are left untouched. The overwrite is recorded as an Iceberg Overwrite operation in the snapshot summary.
-- Unpartitioned: full table replace
INSERT OVERWRITE summary SELECT * FROM staging_summary;
-- Partitioned by region: only the regions in staging are replaced; other regions untouched
INSERT OVERWRITE orders SELECT * FROM staging_orders WHERE region IN ('EU', 'US');
Because the operation is a single atomic Iceberg snapshot, in-flight readers continue to see the pre-overwrite snapshot until the commit lands — there is no intermediate "empty table" window.
Schema Validation
Gnok validates INSERT data at both plan time and execution time:
- Column count — the number of values must match the target column list
- NOT NULL enforcement — all non-nullable columns must be included in the column list (or provided via
SELECT). Missing nullable columns are filled withNULL. - Type casting — safe implicit casts are applied automatically (integer widening, int-to-float). Unsafe casts (e.g., string-to-numeric) are rejected at plan time.
- Precision guard —
INT64values exceeding 2^53 are rejected when casting toFLOAT64, since they cannot be represented exactly.
Writer Options
INSERT writes Parquet data files with these configurable defaults:
| Option | Default | Description |
|---|---|---|
| Target file size | 128 MB | Files are rolled over at this size |
| Compression | zstd(3) | Also supports snappy, gzip, lz4, brotli, none |
| Row group size | 1M rows | Rows per Parquet row group |
| Page size | 1 MB | Target Parquet page size |
| Bloom filters | Enabled | Written per-column for filter pushdown (FPP: 1%) |
| Dictionary encoding | Enabled | For low-cardinality string columns |
| Column statistics | Enabled | Min/max/null counts per column (page-level) |
| Iceberg field IDs | Enabled | Writes PARQUET:field_id metadata for schema evolution |
INSERT Buffer (Staging)
For high-throughput INSERT workloads, Gnok automatically routes small INSERT VALUES statements through a staging buffer that batches multiple INSERTs into fewer, larger Iceberg commits.
How it works:
- INSERTs with up to 5,000 rows are staged to a local RocksDB backend (~0.1ms latency)
- The client receives an immediate acknowledgment — data is durable in the staging backend
- A background flush loop commits staged data to Iceberg every 5 seconds or when 64 MB accumulates
- Multiple small INSERTs become a single Parquet file + single Iceberg snapshot
This reduces snapshot proliferation and produces well-sized Parquet files, dramatically improving both write throughput and subsequent read performance.
Fast-path parser: INSERT VALUES statements that hit the staging buffer bypass the SQL parser entirely on repeat calls. A specialized SIMD-accelerated parser converts VALUES directly into Arrow column arrays with zero intermediate allocations. This achieves 100,000–300,000 rows/sec depending on batch size and column count.
| Setting | Engine reference default | Description |
|---|---|---|
| Buffer threshold | 5,000 rows | INSERT VALUES up to this size use the staging buffer |
| Flush interval | 2s | How often the flush loop checks for ready data |
| Max staleness | 5s | Maximum age before forced flush |
| Target file size | 64 MB | Target Parquet file size per flush |
| Max staging size | 4 GB | Backpressure limit on local staging |
| Max flush retries | 10 | Failures before data moves to dead-letter |
For maximum throughput, use batch sizes of 500 rows per INSERT. This provides the best balance of per-request parsing cost vs. network round-trip overhead, achieving 200K+ rows/sec at high concurrency.
Read-your-writes consistency
The staging buffer does not break SQL's read-your-writes contract. Every read path that could otherwise miss in-flight INSERTs synchronously flushes the affected tables before planning:
SELECTagainst a table with staged rows.UPDATE/DELETE/MERGEover the same target.- Scalar /
EXISTSsubqueries that scan a staged table. SHOW SNAPSHOTSand<table>$snapshotsreads.ALTER TABLE … ADD COLUMN/RENAME COLUMN/DROP COLUMN— flushes under the old schema before the change applies, so staged rows are not stranded.
The synchronous flush runs through the same per-table commit lock as the background loop, so it also serialises against any in-flight background flush — once it returns, every prior INSERT from your session is visible at the Iceberg layer.
The sync flush commits one Iceberg snapshot per staged INSERT statement (rather than coalescing many INSERTs into a single commit, as the background loop does for throughput). This preserves the Iceberg semantic that each INSERT is its own snapshot — needed for AT(TIMESTAMP) and AT(SNAPSHOT) time-travel queries to step through commit history.
You can opt out of the barrier per-session for analytics workloads that prefer lower latency over strong consistency:
SET read_consistency = 'committed'; -- skip the RYW flush
SET read_consistency = 'strong'; -- default; flush before reads
Service defaults are managed by Gnok; use the supported session setting above.
SHOW INSERT BUFFER STATUS
Monitor the staging buffer in real time:
SHOW INSERT BUFFER STATUS;
Returns:
| Column | Type | Description |
|---|---|---|
table_name | VARCHAR | Qualified table name with pending data |
pending_objects | BIGINT | Number of staged objects awaiting flush |
pending_bytes | BIGINT | Total bytes in staging |
oldest_staged_ms | BIGINT | Age of the oldest staged object (milliseconds) |
flush_failures | BIGINT | Consecutive flush failure count (0 = healthy) |
status | VARCHAR | idle, active, backoff, or dead_lettered |
backend | VARCHAR | Staging backend type (rocksdb or postgres) |
Example output:
table_name | pending_objects | pending_bytes | oldest_staged_ms | flush_failures | status | backend
--------------------------------+-----------------+---------------+------------------+----------------+--------+---------
tpcds_catalog.tpcds.orders | 12 | 184320 | 3200 | 0 | active | rocksdb
tpcds_catalog.tpcds.line_items | 3 | 45000 | 800 | 0 | active | rocksdb
Dead-Letter Queue
If the flush loop fails to commit data for a table after max_flush_retries consecutive attempts (e.g., the table was dropped), staged data is moved to a dead-letter store. This prevents poison-pill loops where a failing table blocks other tables from flushing.
Dead-lettered data is preserved in RocksDB for manual inspection. The SHOW INSERT BUFFER STATUS command reports these entries with status dead_lettered.
INSERT SELECT and Bulk Inserts
INSERT SELECT and INSERT VALUES exceeding 5,000 rows use the direct DML path, which dispatches Parquet write tasks to workers via gRPC. This path is designed for bulk operations where the data source is a query or a very large literal batch — not for high-frequency streaming inserts.
The staging buffer (batch sizes up to 5,000 rows) is faster and produces better files than the direct path for INSERT VALUES. Do not increase batch size beyond 5,000 thinking it will be faster — it won't. The staging buffer batches many small INSERTs into optimally-sized Parquet files with fewer Iceberg snapshots, while the direct path creates one snapshot per INSERT.
Use the staging buffer path for streaming/OLTP ingestion. Use INSERT SELECT or COPY INTO for bulk data movement.
The direct path includes safeguards for concurrent workloads:
- Admission control: A semaphore limits concurrent DML operations (default: 64) to prevent coordinator OOM under burst load
- Credential caching: Vended S3 credentials are cached per-table with TTL, eliminating redundant STS calls
- Commit coalescing: Concurrent INSERTs to the same table can have their data files committed together, reducing snapshot count
For partitioned tables, rows are grouped by partition key and routed to partition-owning workers, reducing the small-file problem. For non-partitioned tables, rows are batched and round-robin distributed across workers.
INSERT SELECT always executes distributedly — the source query runs across workers and data is written directly to S3 without flowing through the coordinator.
All inserts use idempotent file naming — file names are derived from a stable job_id plus partition key hash, so retried operations produce identical filenames. Before committing, Gnok checks whether a snapshot with the same job_id already exists, making retries safe against duplicate data.
Output
rows_affected
--------------
3
Zero-row INSERTs are no-ops — no files are written and no Iceberg commit occurs.
Add a RETURNING clause to get the inserted rows back instead of a count.
INSERT ... ON CONFLICT DO UPDATE
Gnok supports Postgres-style upsert syntax for INSERT VALUES:
INSERT INTO orders (order_id, status, updated_at)
VALUES (42, 'shipped', now())
ON CONFLICT (order_id) DO UPDATE
SET status = EXCLUDED.status, updated_at = EXCLUDED.updated_at;
When a row with a conflicting key already exists, the DO UPDATE SET
clause runs against that row using the EXCLUDED pseudo-table to
reference the proposed new row. When no row with the key exists, the
INSERT proceeds normally.
How it works
Under the hood, Gnok rewrites the upsert into a semantically-equivalent
MERGE statement at bind time:
-- Original upsert:
INSERT INTO orders (order_id, status) VALUES (42, 'shipped')
ON CONFLICT (order_id) DO UPDATE SET status = EXCLUDED.status;
-- Effective rewrite:
MERGE INTO orders AS t
USING (VALUES (42, 'shipped')) AS excluded(order_id, status)
ON t.order_id = excluded.order_id
WHEN MATCHED THEN UPDATE SET status = excluded.status
WHEN NOT MATCHED THEN INSERT (order_id, status) VALUES (excluded.order_id, excluded.status);
The rewrite happens in the binder, so the rest of the query engine sees
a regular MERGE and uses the existing merge-on-read execution path.
This means upsert has the same performance characteristics as MERGE:
full target scan, hash-probe, position-delete for updates, new data
file for inserts, single atomic snapshot commit.
Conflict target
The conflict target must be a list of columns, typically the primary key columns:
-- Single-column conflict target
INSERT INTO orders (order_id, status) VALUES (1, 'pending')
ON CONFLICT (order_id) DO UPDATE SET status = EXCLUDED.status;
-- Composite conflict target
INSERT INTO order_lines (order_id, line_id, qty) VALUES (10, 1, 5)
ON CONFLICT (order_id, line_id) DO UPDATE SET qty = EXCLUDED.qty;
Multi-column conflict targets produce an AND-joined join predicate:
t.order_id = excluded.order_id AND t.line_id = excluded.line_id.
EXCLUDED pseudo-table
In the DO UPDATE SET clause, EXCLUDED.col refers to the column value
from the new row (the one that would have been inserted). This lets
you write update expressions that combine existing and proposed data:
-- Bump version counter using the proposed value
INSERT INTO documents (id, content, version)
VALUES (1, 'new content', 2)
ON CONFLICT (id) DO UPDATE
SET content = EXCLUDED.content,
version = EXCLUDED.version + 1;
-- Take max of old vs new
INSERT INTO counters (key, value) VALUES ('visits', 100)
ON CONFLICT (key) DO UPDATE
SET value = GREATEST(value, EXCLUDED.value);
Unqualified column references in SET resolve to the target table
(e.g., value in the last example above).
Gnok's alias resolution for EXCLUDED is case-insensitive —
EXCLUDED.col and excluded.col both work. The same case-insensitive
matching applies to target-table aliases in MERGE statements.
Supported and unsupported forms
| Feature | Supported |
|---|---|
ON CONFLICT (col1, col2) DO UPDATE SET ... | ✅ |
ON CONFLICT ... DO UPDATE SET col = EXCLUDED.col | ✅ |
ON CONFLICT ... DO UPDATE SET col = EXCLUDED.col + literal | ✅ |
INSERT INTO ... VALUES source | ✅ |
INSERT INTO ... SELECT source | ❌ Use MERGE directly |
ON CONFLICT DO NOTHING | ❌ Not yet implemented — use MERGE |
ON CONFLICT ON CONSTRAINT <name> | ❌ Use explicit column list |
ON CONFLICT (col) DO UPDATE SET col = expr WHERE cond | ❌ WHERE on DO UPDATE not yet supported |
Bare ON CONFLICT DO UPDATE (no column list, relies on primary key) | ❌ Must list conflict columns explicitly |
Unsupported forms return a clear error pointing at the limitation — they do not silently fall through.
Atomicity and snapshot semantics
Each upsert statement produces a single Iceberg snapshot containing both the matched-row position deletes and the unmatched-row data file inserts. Concurrent upserts to different rows are serialized at commit time via Iceberg's optimistic concurrency control.
UPDATE
UPDATE orders
SET status = 'cancelled'
WHERE order_id = 42;
UPDATE orders
SET amount = amount * 1.1
WHERE region = 'EU';
How UPDATE Works (Merge-on-Read)
UPDATE uses a Merge-on-Read (MOR) strategy:
- Scan all target data files and evaluate the
WHEREpredicate row-by-row - For each matching row, write a position delete (marking the old row as deleted)
- Apply the
SETassignments and write the updated row as a new data file - Commit both delete files and data files in a single atomic Iceberg snapshot
Rows that were already deleted by prior operations are automatically skipped — both position delete files (the v2 default for partitioned DELETE) and equality delete files (written by DELETE FROM <unpartitioned> WHERE col = N and similar predicate-driven deletes on identifier-only tables). Earlier the UPDATE scan loaded only position deletes, so a row that had been deleted via an equality predicate would resurrect with the new column value applied — a DELETE FROM t WHERE id = 5 followed by UPDATE t SET v = v + 1 would re-introduce row id = 5 with the bumped value. The equality-delete-file list is now propagated through the coordinator → scan-task → executor pipeline and applied per-batch.
This means UPDATEs are cheap for small changes but reads must merge delete files at query time. Run OPTIMIZE PURGE DELETES periodically to consolidate.
UPDATE currently scans all data files in the table — the WHERE predicate is not used for file-level or partition-level pruning. For large tables, ensure your UPDATE has a selective WHERE clause.
Supported Assignment Expressions
SET clauses support:
-- Literals
UPDATE t SET col = 42;
UPDATE t SET col = 'string';
UPDATE t SET col = NULL;
-- Column references
UPDATE t SET col1 = col2;
-- Arithmetic (+, -, *, /, %, ||)
UPDATE t SET price = price * 1.1;
UPDATE t SET count = count - 1;
UPDATE t SET name = first_name || ' ' || last_name;
-- Functions
UPDATE t SET value = ROUND(value, 2);
UPDATE t SET name = UPPER(name);
UPDATE t SET col = COALESCE(col, 'default');
UPDATE t SET col = ABS(col);
-- CAST
UPDATE t SET col = CAST(value AS INT);
Gnok automatically casts the expression result to the target column type when safe. Numeric types are promoted to the wider type before evaluation (e.g., INT32 + FLOAT64 promotes to FLOAT64).
Subqueries, CASE expressions, and window functions are not supported in SET clauses. Multi-table UPDATE ... FROM / UPDATE ... JOIN is not supported — use MERGE instead.
UPDATE ... FROM <source> is rejected directionally
The PG-style UPDATE target SET ... FROM source WHERE ... shape
is rejected at bind time with a directional message pointing at
MERGE INTO:
UPDATE ... FROM <source> is not supported — use
MERGE INTO target USING source ON ... WHEN MATCHED THEN UPDATE SET ...
Rewriting an UPDATE ... FROM as MERGE is mechanical:
-- Reject
UPDATE orders o
SET status = s.new_status
FROM staged_status_updates s
WHERE o.order_id = s.order_id;
-- Equivalent MERGE
MERGE INTO orders o
USING staged_status_updates s
ON o.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET status = s.new_status;
DECIMAL(38, _) overflow during SET
UPDATE ... SET col = col + <value> where col is
DECIMAL(38, s) may raise an overflow error if the addition would
exceed 10^38 - 1 — see Decimal arithmetic
overflow in the
data-types reference. The same applies to expressions inside
INSERT ... VALUES and MERGE ... WHEN MATCHED THEN UPDATE SET.
Parallel Execution
For tables with 4+ data files, Gnok parallelizes the UPDATE across workers. Files are grouped by partition for locality — position delete files are colocated with their corresponding data files.
| Setting | Default |
|---|---|
| Min files for parallel | 4 |
| Max files per task | 25 |
| Max tasks per partition | 4 |
| Memory budget per task | 1 GB |
UPDATE Without WHERE
Omitting the WHERE clause updates every row in the table:
UPDATE orders SET status = 'pending';
Output
rows_affected
--------------
42
Add a RETURNING clause to get the updated rows (post-image) back instead of a count.
DELETE
DELETE FROM orders WHERE order_date < '2023-01-01';
DELETE FROM orders WHERE status = 'cancelled';
-- Delete all rows
DELETE FROM orders;
How DELETE Works
DELETE uses a Merge-on-Read position delete strategy (Iceberg v2). For each data file, the WHERE predicate is evaluated row-by-row and matching row positions are recorded in Parquet delete files with schema (file_path STRING, pos INT64).
Position deletes are resolved at read time by the scan operator. Accumulated delete files should be periodically consolidated with OPTIMIZE PURGE DELETES.
The read-path delete index uses a Vec<i64> with binary search for files with fewer than 50,000 deleted positions, and a RoaringBitmap for larger sets (service-managed threshold).
Deletion Vectors (v3)
For tables using Iceberg format version 3, Gnok can read deletion vectors — a compact representation using 64-bit Roaring bitmaps stored in Puffin files. Deletion vectors are matched to data files by referenced_data_file path and merged with any existing position deletes during scan execution. DV write support is planned for a future release.
WHERE col IN (SELECT ...) — uncorrelated subqueries
DELETE and UPDATE accept an uncorrelated subquery on the
right-hand side of IN. The coordinator evaluates the subquery
once, materialises the result as a literal IN-list, and sends
that list to workers as part of the predicate:
-- Delete every order whose customer is in a staging table
DELETE FROM orders
WHERE customer_id IN (SELECT id FROM staged_cancellations);
-- Update only the rows whose key is in a filtered set
UPDATE orders SET status = 'archived'
WHERE order_id IN (SELECT order_id FROM legacy_orders WHERE archived);
-- NOT IN works the same way
DELETE FROM orders
WHERE customer_id NOT IN (SELECT id FROM active_customers);
Because the subquery is materialised before the worker dispatch, the
inner query can reference any tables / joins / aggregates — the only
constraint is that it must produce a single column whose type is
implicitly castable to the outer column. A subquery that returns
millions of keys is fine; the planner streams them into a large IN
list rather than rebuilding it per row.
Correlated subqueries in DML WHERE (WHERE col IN (SELECT … WHERE inner.x = outer.y))
are rejected at bind time — per-row subquery evaluation isn't wired
through the DML path. Rewrite as MERGE:
-- Reject (correlated):
DELETE FROM orders o
WHERE EXISTS (SELECT 1 FROM staged s WHERE s.order_id = o.id);
-- Rewrite:
MERGE INTO orders o
USING staged s ON s.order_id = o.id
WHEN MATCHED THEN DELETE;
DELETE FROM <target> USING <source> is rejected directionally
The PG-style multi-table DELETE FROM t USING u WHERE t.id = u.id
shape parses but the column-resolution path does not see the
USING-side alias, so u.col references would previously surface
as Column 'u' not found. This shape is now rejected at bind time
with a directional message pointing at the supported workarounds:
-- Reject:
DELETE FROM orders USING staged_cancellations s WHERE orders.id = s.id;
-- Rewrite as IN-subquery (the source is a single column):
DELETE FROM orders WHERE id IN (SELECT id FROM staged_cancellations);
-- Or as MERGE (the source has multiple useful columns):
MERGE INTO orders t USING staged_cancellations s
ON t.id = s.id
WHEN MATCHED THEN DELETE;
Partition Pruning
DELETE uses the same scan planner as SELECT queries, so partition pruning is applied when the WHERE clause references partition columns.
Parallel Execution
For tables with 4+ data files, DELETE is parallelized across workers. Files are grouped by partition to ensure delete files are colocated with their data files.
| Setting | Default |
|---|---|
| Min files for parallel | 4 |
| Max files per task | 50 |
| Max tasks per partition | 4 |
| Memory budget per task | 512 MB |
TRUNCATE TABLE
To delete all rows without scanning, use TRUNCATE:
TRUNCATE TABLE orders;
TRUNCATE is a metadata-only operation — it creates a new empty Iceberg snapshot without scanning or deleting any data files from object storage. Data files are physically removed by a subsequent VACUUM operation. This is significantly faster than DELETE without a WHERE clause.
Output
rows_affected
--------------
1200
Add a RETURNING clause to get the deleted rows back instead of a count.
RETURNING
INSERT, UPDATE, and DELETE accept an optional RETURNING clause that turns the statement into a query: instead of a rows_affected count, the statement returns a result set built from the affected rows. This matches PostgreSQL's RETURNING and works over both the simple and extended (prepared-statement) wire protocols.
-- All columns of the affected rows
DELETE FROM orders WHERE status = 'cancelled' RETURNING *;
-- Specific columns
INSERT INTO orders (order_id, amount) VALUES (1001, 250.00) RETURNING order_id;
-- Expressions and aliases
UPDATE orders SET amount = amount * 1.1 WHERE region = 'us'
RETURNING order_id, amount AS new_amount;
The RETURNING list accepts * (all columns of the target table), bare column names, and arbitrary expressions over the target table's columns, each optionally aliased with AS. An unaliased expression is named ?column?, following PostgreSQL.
Which row image is returned:
| Statement | Rows returned |
|---|---|
INSERT | The values being inserted, projected from the VALUES / SELECT source. |
UPDATE | The post-image — column values after the SET assignments are applied. |
DELETE | The pre-image — the rows as they existed before deletion. |
-- DELETE ... RETURNING *
order_id | status | amount
---------+-----------+--------
2007 | cancelled | 120.00
2011 | cancelled | 64.50
INSERT ... RETURNING projects the values from the statement's own source rows. It does not read back server-generated values — column DEFAULTs, sequence / IDENTITY values, and computed columns are not reflected in the returned rows. If you need a server-assigned value (e.g. an auto-generated key), run a follow-up SELECT.
Limitations:
MERGE ... RETURNINGis not supported.INSERT ... DEFAULT VALUES ... RETURNINGis not supported (there is no explicit source to project from).- Expressions in the
RETURNINGlist resolve only against the target table's columns.
MERGE
MERGE combines INSERT, UPDATE, and DELETE into a single atomic statement. It matches rows between a source and target using a join condition, then applies different actions depending on whether a match is found.
Syntax
MERGE INTO target_table [AS alias]
USING source [AS alias]
ON join_condition
[WHEN MATCHED [AND condition] THEN UPDATE SET col = expr, ... | DELETE]+
[WHEN NOT MATCHED [AND condition] THEN INSERT (cols) VALUES (vals)]
[WHEN NOT MATCHED BY SOURCE [AND condition] THEN UPDATE SET col = expr, ... | DELETE]+;
Examples
-- Upsert: update existing rows, insert new ones
MERGE INTO customers AS t
USING staging_customers AS s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET t.name = s.name, t.email = s.email, t.updated_at = now()
WHEN NOT MATCHED THEN INSERT (id, name, email, updated_at) VALUES (s.id, s.name, s.email, now());
-- Conditional update with delete
MERGE INTO target AS t
USING source AS s
ON t.id = s.id
WHEN MATCHED AND s.deleted = true THEN DELETE
WHEN MATCHED THEN UPDATE SET t.value = s.value, t.updated_at = now()
WHEN NOT MATCHED THEN INSERT (id, value, updated_at) VALUES (s.id, s.value, now())
WHEN NOT MATCHED BY SOURCE THEN DELETE;
-- Source can be a subquery
MERGE INTO inventory AS t
USING (SELECT product_id, SUM(quantity) AS qty FROM shipments GROUP BY product_id) AS s
ON t.product_id = s.product_id
WHEN MATCHED THEN UPDATE SET t.quantity = t.quantity + s.qty
WHEN NOT MATCHED THEN INSERT (product_id, quantity) VALUES (s.product_id, s.qty);
Source Types
The USING source can be a table, a subquery, or a VALUES clause:
-- Subquery
MERGE INTO target AS t
USING (SELECT id, value FROM staging WHERE active = true) AS s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET t.value = s.value;
-- VALUES clause
MERGE INTO target AS t
USING (VALUES (1, 'a'), (2, 'b')) AS s(id, name)
ON t.id = s.id
WHEN NOT MATCHED THEN INSERT (id, name) VALUES (s.id, s.name);
CTEs (WITH ... AS) are also supported as the outer query context.
Clause Evaluation
- Multiple
WHEN MATCHEDclauses are supported — they are evaluated in order, and the first match wins - One
WHEN NOT MATCHEDclause is allowed for inserting unmatched source rows - Multiple
WHEN NOT MATCHED BY SOURCEclauses are supported — evaluated in order - Each clause can have an optional
AND conditionto add further filtering - NULL join keys never match (SQL standard:
NULL != NULL). Source rows with NULL keys are treated as unmatched and are eligible forWHEN NOT MATCHEDinsertion.
ON Condition Requirements
The ON condition must be an equi-join using direct column references:
-- Supported
ON t.id = s.id
ON t.id = s.id AND t.region = s.region
-- Not supported
ON t.id = s.id OR t.name = s.name -- No OR
ON upper(t.id) = s.id -- No function calls
Supported SET expressions in WHEN MATCHED UPDATE
The SET clause in a WHEN MATCHED THEN UPDATE action supports
literals, target and source column references, and arithmetic/comparison
operators. Source column references use either the source alias or the
EXCLUDED keyword (for ON CONFLICT rewrites):
WHEN MATCHED THEN UPDATE
SET version = s.version,
last_seen = now(),
score = target.score + s.delta,
status = CASE WHEN s.active THEN 'on' ELSE 'off' END
Alias matching is case-insensitive — T.col, t.col, S.col, and
s.col all resolve correctly regardless of how they were declared in
the USING ... AS clause.
Mixed-width integer arithmetic is auto-promoted. An expression like
int32_col + 10 (where the literal 10 is parsed as Int64) is
widened to the wider integer type before evaluation. All Int8 /
Int16 / Int32 / Int64 pairs are handled. Without this, mixed-
type arithmetic would fail with a Type mismatch in binary operation
error.
Cardinality enforcement
A MERGE source must produce at most one row per matched target
row. If two or more source rows match the same target row under
the ON condition, the statement aborts with a clear error rather
than silently applying the actions in arbitrary order:
MERGE source produced N source rows for the same target row;
the cardinality is undefined. Pre-dedup the source on the ON-condition
columns, or add a more selective ON predicate.
Earlier the engine quietly picked one of the matching source rows
and applied that row's WHEN MATCHED action, which depended on
worker / file ordering. Matches the SQL-standard semantics.
How MERGE Works Internally
- Build phase: The source is materialized into an in-memory hash table using a collision-free byte encoding of join keys. The hash table is bounded by configurable limits (default: 1M rows / 512 MB). Exceeding these limits returns an error — there is no spill to disk.
- Probe phase: All target files are scanned row-by-row, each row probed against the hash table.
WHEN NOT MATCHED BY SOURCEclauses are evaluated for unmatched target rows during this phase. - Insert phase: Source rows not found in the target (unmatched keys) are collected and passed through the
WHEN NOT MATCHEDclause. - Write phase: Position deletes for modified/deleted rows and new data files for updated/inserted rows are written in batches.
- All files are committed in a single atomic Iceberg snapshot.
MERGE always performs a full scan of the target table — the ON condition is not used for partition or file pruning.
Limits
| Setting | Default | Description |
|---|---|---|
| Maximum source rows | 1,000,000 | Maximum source rows materialized into the hash table |
| Maximum source memory (MB) | 512 | Maximum source memory (Arrow in-memory size) |
| Min files for parallel | 4 | Below this, MERGE runs as a single task |
| Max files per task | 20 | Target files per parallel task |
| Memory budget per task | 2 GB | Per-task memory limit |
Output
MERGE returns a summary of the actions taken:
rows_inserted | rows_updated | rows_deleted
--------------+--------------+-------------
150 | 320 | 45
COPY INTO
Bulk load data from stages into Iceberg tables. COPY INTO is optimized for high-throughput ingestion with SIMD-accelerated CSV parsing, adaptive file splitting, and distributed shuffle for partitioned tables.
Stages
Gnok uses stages as managed storage locations for data files. A stage reference has the form:
@stage/<stage_name>/<run_id>/<filename>
@stage/<stage_name>/<run_id>/ -- directory (all files)
Stages are implicit namespaces backed by a Gnok-managed S3 bucket — there is no CREATE STAGE DDL. Any stage name can be used when uploading files.
Upload Workflow
Before running COPY INTO, upload data files to a stage via presigned URLs:
- Request a presigned URL:
POST /v1/stages/presign
Authorization: Bearer <jwt>
{
"stage": "uploads",
"runId": "batch_001",
"filename": "data.parquet",
"contentType": "application/octet-stream"
}
Response:
{
"method": "PUT",
"url": "https://s3.amazonaws.com/bucket/.../data.parquet?X-Amz-Signature=...",
"stageRef": "@stage/uploads/batch_001/data.parquet",
"expiresInSecs": 900
}
-
Upload the file with a direct
PUTto the presigned URL. -
Run COPY INTO referencing the stage path.
For large files, use the multipart upload endpoints:
| Endpoint | Purpose |
|---|---|
POST /v1/stages/presign/multipart/init | Initialize multipart upload, get per-part presigned URLs |
POST /v1/stages/presign/multipart/complete | Finalize after all parts are uploaded |
POST /v1/stages/presign/multipart/abort | Cancel a multipart upload |
Credentials are automatically vended by the catalog service using the JWT token — no manual S3 configuration is required.
Basic Usage
-- Load Parquet files from a stage
COPY INTO orders
FROM '@stage/uploads/batch_001/'
FILE_FORMAT = (TYPE = PARQUET);
-- Load CSV with options
COPY INTO orders
FROM '@stage/raw/run_20240115/'
FILE_FORMAT = (
TYPE = CSV,
DELIMITER = ',',
QUOTE = '"',
SKIP_ROWS = 1,
NULL_STRING = 'NULL'
);
-- Load newline-delimited JSON
COPY INTO events
FROM '@stage/events/daily_2025/'
FILE_FORMAT = (TYPE = JSON);
-- Load a single file
COPY INTO orders
FROM '@stage/uploads/batch_001/data.parquet'
FILE_FORMAT = (TYPE = PARQUET);
PATTERN
Filter files by regex pattern:
-- Only load .parquet files matching a pattern
COPY INTO orders
FROM '@stage/uploads/batch_001/'
FILE_FORMAT = (TYPE = PARQUET)
PATTERN = '.*orders.*\.parquet';
INFER_SCHEMA
Before loading, you can inspect the schema of staged files:
SELECT * FROM INFER_SCHEMA(
LOCATION => '@stage/uploads/run001/',
FILE_FORMAT => 'CSV',
OPTIONS => (
SAMPLE_SIZE => 1000,
MAX_FILES => 10,
NULL_THRESHOLD => 0.95,
INT_TO_BIGINT => TRUE,
FLOAT_TO_DOUBLE => TRUE
)
);
Returns: column_name, inferred_type, nullable, confidence.
| Option | Default | Description |
|---|---|---|
SAMPLE_SIZE | 1000 | Rows sampled per file (CSV/JSON only; Parquet reads metadata directly) |
MAX_FILES | 10 | Maximum files to read for inference |
NULL_THRESHOLD | 0.95 | Column marked nullable if null ratio exceeds this |
INT_TO_BIGINT | true | Promote INT to BIGINT |
FLOAT_TO_DOUBLE | true | Promote FLOAT to DOUBLE |
FILE_FORMAT Options
CSV:
| Option | Default | Description |
|---|---|---|
DELIMITER | , | Field separator character |
QUOTE | " | Quote character for fields containing delimiters |
ESCAPE | "" (doubled) | Escape character (default is RFC 4180 quote doubling) |
SKIP_ROWS | 0 | Number of header rows to skip |
NULL_STRING | — | String value to interpret as NULL |
CSV parsing uses a SIMD-accelerated parser (AVX2 on x86_64, NEON on ARM) that processes structural characters in parallel, achieving 500-700 MB/s throughput.
Parquet and JSON use default reader settings with no additional format options.
Error Handling
| Option | Default | Description |
|---|---|---|
ON_ERROR | ABORT | ABORT (fail on first error), CONTINUE (skip bad rows), or SKIP_FILE (skip entire file on error) |
MAX_ERRORS | 1000 | Max errors before aborting in CONTINUE mode |
COPY INTO orders
FROM '@stage/uploads/batch_001/'
FILE_FORMAT = (TYPE = CSV, SKIP_ROWS = 1)
ON_ERROR = CONTINUE;
Distributed Execution
Gnok automatically parallelizes COPY INTO across workers with adaptive file planning:
| File Size | Strategy |
|---|---|
| < 32 MB | Batched — small files grouped together (up to 32 files / 256 MB per batch) |
| 32-256 MB | Single task — one file per task |
| > 256 MB | Split — file split into ~64 MB byte ranges processed in parallel |
For partitioned tables, Gnok routes rows to partition-owning workers via a shuffle service (gRPC-based, with backpressure). This ensures each partition gets optimally-sized data files rather than many small files.
During execution, the coordinator vends temporary STS credentials to workers — workers never access stages directly with long-lived credentials.
Limits
| Limit | Default | Description |
|---|---|---|
| Maximum files | 10,000 | Maximum files per COPY INTO |
| Maximum source size (GB) | 100 | Maximum source data size in GB |
Output
COPY INTO returns a summary:
files_loaded | rows_loaded | bytes_loaded
-------------+-------------+-------------
12 | 2500000 | 1073741824
REGISTER FILES INTO
Register existing Parquet files in object storage as part of an Iceberg table without copying or rewriting data. Gnok reads only the Parquet footers (metadata) and commits the file references as a new Iceberg snapshot.
This is useful for:
- Migrating existing Parquet data lakes into Iceberg with zero data movement
- Registering files produced by external systems (Spark, Flink, etc.)
- Bulk onboarding historical data already in S3
Syntax
REGISTER FILES INTO [catalog.][schema.]table
FROM 'source_uri'
[FILE_FORMAT PARQUET];
FILE_FORMAT defaults to PARQUET and is currently the only supported format.
Examples
-- Register all Parquet files under an S3 prefix
REGISTER FILES INTO orders FROM 's3://data-lake/warehouse/orders/';
-- With explicit schema and format
REGISTER FILES INTO analytics.events FROM 's3://data-lake/events/2025/' FILE_FORMAT PARQUET;
-- Three-part table name
REGISTER FILES INTO prod.public.orders FROM 's3://data-lake/warehouse/orders/';
-- s3a:// URIs are also accepted
REGISTER FILES INTO orders FROM 's3a://data-lake/warehouse/orders/';
How It Works
- List files — discovers all
*.parquetfiles at the source URI (sorted by path) - Validate schema — reads the first file's Parquet footer and checks that all non-nullable table columns exist in the Parquet schema. Extra Parquet columns not in the table schema are ignored with a warning.
- Read footers concurrently — reads Parquet metadata (up to 32 files in parallel) to extract row counts, file sizes, and null value counts. No row data is transferred.
- Commit — registers all files as a single atomic Iceberg snapshot using the same commit path as INSERT and COPY INTO
Output
files_registered | rows_registered | bytes_registered
-----------------+-----------------+-----------------
48 | 12500000 | 2147483648
REGISTER FILES vs COPY INTO
| REGISTER FILES INTO | COPY INTO | |
|---|---|---|
| Data transfer | None — reads Parquet footers only | Reads and rewrites all data to table location |
| File location | Files stay at their original S3 path | Files are written to the table's managed location |
| Supported formats | Parquet only | Parquet, CSV, JSON |
| Partitioned tables | Not supported (files registered as unpartitioned) | Supported with automatic partition routing |
| Column statistics | Null counts only (min/max bounds not populated) | Full statistics written during ingest |
| Source | s3:// or s3a:// URIs | @stage/... references |
| Error handling | All-or-nothing | Configurable (ON_ERROR, MAX_ERRORS) |
Requirements
- The target table must already have at least one snapshot (i.e., at least one prior INSERT or COPY INTO). Registering into a brand-new empty table is not yet supported.
- Source files must be Parquet format.
- The Parquet schema must contain all non-nullable columns defined in the table schema.
Transactions
Gnok supports multi-statement ACID transactions. Outside an explicit
transaction the engine runs in autocommit mode — each INSERT, UPDATE,
DELETE, MERGE, and COPY INTO is its own atomic Iceberg snapshot commit. Inside
a BEGIN … COMMIT block, all of the block's DML is staged and committed as a
single atomic unit, and ROLLBACK discards it.
Transactions work identically over all protocols — the PostgreSQL wire
protocol, Flight SQL, and the HTTP API. (Flight SQL and HTTP key the
transaction to the x-gnok-session-id header, which must be a stable UUID per
connection.)
Transaction control statements
| Statement | Aliases | Effect |
|---|---|---|
BEGIN | BEGIN TRANSACTION, START TRANSACTION | Start a transaction. An optional ISOLATION LEVEL … clause sets the level. |
COMMIT | END | Atomically commit all staged changes. |
ROLLBACK | ROLLBACK TRANSACTION, ABORT | Discard all staged changes (and delete any data files written by the transaction). |
SET TRANSACTION ISOLATION LEVEL … | Set the isolation level for the current transaction. | |
SHOW TRANSACTIONS | List the active transactions on the coordinator. |
BEGIN;
INSERT INTO accounts VALUES (1, 100);
UPDATE accounts SET balance = balance - 50 WHERE id = 1;
INSERT INTO ledger VALUES (1, -50);
COMMIT; -- all three land in one atomic commit, or none do
A bare BEGIN/COMMIT/ROLLBACK with no DML in between is a clean no-op and
returns a status message, so connection pools and ORMs that bracket every query
in a transaction shape work unchanged. The Iceberg-flavoured
ROLLBACK TABLE <name> TO SNAPSHOT … form is a separate snapshot DDL operation,
not transaction control (see the DDL reference).
Atomicity
The block commits as one unit:
- Single table — one Iceberg snapshot.
- Multiple tables in the same catalog — one catalog
commit_transaction(all tables' updates in a single multi-table commit; a conflict on any table aborts the whole transaction). - Multiple catalogs — committed in one underlying database transaction via the catalog's
/v1/transactions/commitendpoint. (Cross-tenant transactions are rejected.)
ROLLBACK (or a failed COMMIT) discards the staged changes and eagerly
deletes any data files the transaction already wrote to object storage, so an
aborted transaction leaves nothing behind.
Isolation levels
| Level | Read behavior |
|---|---|
READ COMMITTED (default) | Each statement reads the latest committed snapshot at the time it runs. |
REPEATABLE READ | The snapshot is pinned at the transaction's first read; all later reads see that same snapshot. |
SERIALIZABLE | Snapshot reads as in REPEATABLE READ, plus the read set is validated at commit time — a concurrent write to anything the transaction read aborts the commit (prevents write skew). |
BEGIN ISOLATION LEVEL SERIALIZABLE;
SELECT balance FROM accounts WHERE id = 1; -- read set is tracked
UPDATE accounts SET balance = 0 WHERE id = 1;
COMMIT; -- aborts with a serialization failure if `accounts` changed underneath
Concurrency is optimistic (no locks held during execution). At commit the
catalog asserts the base snapshot is unchanged; conflicts on append/delete-style
statements are automatically re-based and retried with backoff, while
SERIALIZABLE read-set conflicts surface to the client as a serialization
failure to retry. See Optimistic Concurrency Control.
Aborted transactions
If any statement inside a transaction errors, the transaction enters the
aborted state: subsequent statements are rejected with
current transaction is aborted, commands ignored until end of transaction block
(SQLSTATE 25P02) until you issue ROLLBACK, and a COMMIT in this state rolls
back instead of committing. This matches PostgreSQL.
DDL inside a transaction
DDL (CREATE / ALTER / DROP) is not transactional — a DDL statement
inside a BEGIN block auto-commits the transaction's pending DML along with the
DDL itself, and a later ROLLBACK cannot undo it. Keep DDL out of transactions
you intend to roll back.
Savepoints
SAVEPOINT and nested transactions are not supported — they are rejected
with SAVEPOINT / nested transactions are not supported (SQLSTATE 0A000)
rather than silently ignored, so a client relying on partial rollback fails
loudly instead of getting wrong results.
Snapshot Isolation
Under the default READ COMMITTED level (and for every autocommit statement),
each DML statement captures the current snapshot ID at planning time, reads a
consistent view of that snapshot, and commits its new data/delete files as a new
snapshot that builds on the base. REPEATABLE READ and SERIALIZABLE extend
this by pinning the snapshot across the whole transaction (see above).
Optimistic Concurrency Control
Gnok uses optimistic concurrency control (OCC) — no locks are held during execution. At commit time, the catalog verifies that the base snapshot has not changed:
- REST catalog: The commit includes an
assert-ref-snapshot-idrequirement. A409 Conflictresponse triggers a retry. - Glue catalog: The commit uses
UpdateTablewith aversion_idcondition. AConcurrentModificationExceptiontriggers a retry.
If a concurrent writer committed between the plan and commit phases, the operation is retried with exponential backoff:
| Retry | Delay |
|---|---|
| 1st | ~100 ms |
| 2nd | ~200 ms |
| 3rd | ~400 ms |
| 4th | ~800 ms |
Delays include ±25% jitter and are capped at 5 seconds. After 4 retries (5 total attempts), the operation fails with a conflict error.
Idempotency
Each DML operation is assigned a stable job_id that is stored in the Iceberg snapshot summary. Before committing, Gnok checks whether a snapshot with the same job_id already exists. If so, the commit is a no-op — making retries safe against duplicate data.
File names are also deterministic (derived from job_id + partition key hash), so retried writes produce identical files.
Orphan File Cleanup
When a commit definitively fails (conflict, storage error), Gnok attempts to clean up the metadata files it wrote. When the outcome is ambiguous (timeout, catalog error), the files are recorded as orphan candidates for later reconciliation.
VACUUM handles long-term cleanup:
- Expired snapshots and their unreferenced data files are removed
- Orphan files older than 3 days (default) are deleted
- A minimum of 5 snapshots are always retained
Snapshot Pinning
Sessions can pin a table's snapshot to get repeatable reads within a time window:
| Setting | Default | Description |
|---|---|---|
| Snapshot mode | pinned | pinned (cache for TTL), latest (always current), or explicit:<id> (time travel) |
| Snapshot TTL (ms) | 15000 | How long a pinned snapshot is reused before re-resolving |
| Maximum pinned tables | 100 | Maximum tables with pinned snapshots per session |
With pinned mode (default), a session may read from a slightly stale snapshot for up to 15 seconds after a concurrent write commits. Use latest mode if read-your-writes consistency is required.