Skip to main content

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.

note

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 with NULL.
  • 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 — INT64 values exceeding 2^53 are rejected when casting to FLOAT64, since they cannot be represented exactly.

Writer Options​

INSERT writes Parquet data files with these configurable defaults:

OptionDefaultDescription
Target file size128 MBFiles are rolled over at this size
Compressionzstd(3)Also supports snappy, gzip, lz4, brotli, none
Row group size1M rowsRows per Parquet row group
Page size1 MBTarget Parquet page size
Bloom filtersEnabledWritten per-column for filter pushdown (FPP: 1%)
Dictionary encodingEnabledFor low-cardinality string columns
Column statisticsEnabledMin/max/null counts per column (page-level)
Iceberg field IDsEnabledWrites 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:

  1. INSERTs with up to 5,000 rows are staged to a local RocksDB backend (~0.1ms latency)
  2. The client receives an immediate acknowledgment — data is durable in the staging backend
  3. A background flush loop commits staged data to Iceberg every 5 seconds or when 64 MB accumulates
  4. 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.

SettingEngine reference defaultDescription
Buffer threshold5,000 rowsINSERT VALUES up to this size use the staging buffer
Flush interval2sHow often the flush loop checks for ready data
Max staleness5sMaximum age before forced flush
Target file size64 MBTarget Parquet file size per flush
Max staging size4 GBBackpressure limit on local staging
Max flush retries10Failures before data moves to dead-letter
tip

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:

  • SELECT against a table with staged rows.
  • UPDATE / DELETE / MERGE over the same target.
  • Scalar / EXISTS subqueries that scan a staged table.
  • SHOW SNAPSHOTS and <table>$snapshots reads.
  • 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:

ColumnTypeDescription
table_nameVARCHARQualified table name with pending data
pending_objectsBIGINTNumber of staged objects awaiting flush
pending_bytesBIGINTTotal bytes in staging
oldest_staged_msBIGINTAge of the oldest staged object (milliseconds)
flush_failuresBIGINTConsecutive flush failure count (0 = healthy)
statusVARCHARidle, active, backoff, or dead_lettered
backendVARCHARStaging 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.

caution

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).

note

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​

FeatureSupported
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:

  1. Scan all target data files and evaluate the WHERE predicate row-by-row
  2. For each matching row, write a position delete (marking the old row as deleted)
  3. Apply the SET assignments and write the updated row as a new data file
  4. 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.

note

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.

SettingDefault
Min files for parallel4
Max files per task25
Max tasks per partition4
Memory budget per task1 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.

SettingDefault
Min files for parallel4
Max files per task50
Max tasks per partition4
Memory budget per task512 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:

StatementRows returned
INSERTThe values being inserted, projected from the VALUES / SELECT source.
UPDATEThe post-image — column values after the SET assignments are applied.
DELETEThe 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 reflects source values, not server-computed ones

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 ... RETURNING is not supported.
  • INSERT ... DEFAULT VALUES ... RETURNING is not supported (there is no explicit source to project from).
  • Expressions in the RETURNING list 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 MATCHED clauses are supported — they are evaluated in order, and the first match wins
  • One WHEN NOT MATCHED clause is allowed for inserting unmatched source rows
  • Multiple WHEN NOT MATCHED BY SOURCE clauses are supported — evaluated in order
  • Each clause can have an optional AND condition to add further filtering
  • NULL join keys never match (SQL standard: NULL != NULL). Source rows with NULL keys are treated as unmatched and are eligible for WHEN NOT MATCHED insertion.

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​

  1. 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.
  2. Probe phase: All target files are scanned row-by-row, each row probed against the hash table. WHEN NOT MATCHED BY SOURCE clauses are evaluated for unmatched target rows during this phase.
  3. Insert phase: Source rows not found in the target (unmatched keys) are collected and passed through the WHEN NOT MATCHED clause.
  4. Write phase: Position deletes for modified/deleted rows and new data files for updated/inserted rows are written in batches.
  5. All files are committed in a single atomic Iceberg snapshot.
note

MERGE always performs a full scan of the target table — the ON condition is not used for partition or file pruning.

Limits​

SettingDefaultDescription
Maximum source rows1,000,000Maximum source rows materialized into the hash table
Maximum source memory (MB)512Maximum source memory (Arrow in-memory size)
Min files for parallel4Below this, MERGE runs as a single task
Max files per task20Target files per parallel task
Memory budget per task2 GBPer-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:

  1. 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
}
  1. Upload the file with a direct PUT to the presigned URL.

  2. Run COPY INTO referencing the stage path.

For large files, use the multipart upload endpoints:

EndpointPurpose
POST /v1/stages/presign/multipart/initInitialize multipart upload, get per-part presigned URLs
POST /v1/stages/presign/multipart/completeFinalize after all parts are uploaded
POST /v1/stages/presign/multipart/abortCancel 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.

OptionDefaultDescription
SAMPLE_SIZE1000Rows sampled per file (CSV/JSON only; Parquet reads metadata directly)
MAX_FILES10Maximum files to read for inference
NULL_THRESHOLD0.95Column marked nullable if null ratio exceeds this
INT_TO_BIGINTtruePromote INT to BIGINT
FLOAT_TO_DOUBLEtruePromote FLOAT to DOUBLE

FILE_FORMAT Options​

CSV:

OptionDefaultDescription
DELIMITER,Field separator character
QUOTE"Quote character for fields containing delimiters
ESCAPE"" (doubled)Escape character (default is RFC 4180 quote doubling)
SKIP_ROWS0Number 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​

OptionDefaultDescription
ON_ERRORABORTABORT (fail on first error), CONTINUE (skip bad rows), or SKIP_FILE (skip entire file on error)
MAX_ERRORS1000Max 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 SizeStrategy
< 32 MBBatched — small files grouped together (up to 32 files / 256 MB per batch)
32-256 MBSingle task — one file per task
> 256 MBSplit — 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​

LimitDefaultDescription
Maximum files10,000Maximum files per COPY INTO
Maximum source size (GB)100Maximum 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​

  1. List files — discovers all *.parquet files at the source URI (sorted by path)
  2. 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.
  3. 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.
  4. 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 INTOCOPY INTO
Data transferNone — reads Parquet footers onlyReads and rewrites all data to table location
File locationFiles stay at their original S3 pathFiles are written to the table's managed location
Supported formatsParquet onlyParquet, CSV, JSON
Partitioned tablesNot supported (files registered as unpartitioned)Supported with automatic partition routing
Column statisticsNull counts only (min/max bounds not populated)Full statistics written during ingest
Sources3:// or s3a:// URIs@stage/... references
Error handlingAll-or-nothingConfigurable (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​

StatementAliasesEffect
BEGINBEGIN TRANSACTION, START TRANSACTIONStart a transaction. An optional ISOLATION LEVEL … clause sets the level.
COMMITENDAtomically commit all staged changes.
ROLLBACKROLLBACK TRANSACTION, ABORTDiscard 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 TRANSACTIONSList 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/commit endpoint. (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​

LevelRead behavior
READ COMMITTED (default)Each statement reads the latest committed snapshot at the time it runs.
REPEATABLE READThe snapshot is pinned at the transaction's first read; all later reads see that same snapshot.
SERIALIZABLESnapshot 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-id requirement. A 409 Conflict response triggers a retry.
  • Glue catalog: The commit uses UpdateTable with a version_id condition. A ConcurrentModificationException triggers a retry.

If a concurrent writer committed between the plan and commit phases, the operation is retried with exponential backoff:

RetryDelay
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:

SettingDefaultDescription
Snapshot modepinnedpinned (cache for TTL), latest (always current), or explicit:<id> (time travel)
Snapshot TTL (ms)15000How long a pinned snapshot is reused before re-resolving
Maximum pinned tables100Maximum 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.