Streaming Ingestion
Gnok supports real-time streaming ingestion with sub-second latency. Records are buffered in memory and committed as Iceberg snapshots, making them immediately queryable.
Architecture
The streaming ingestion pipeline has four stages:
Producers Gnok Coordinator Object Storage
+---------+ +---------------------+ +-------------+
| gRPC / | ----> | Streaming API | | |
| HTTP / | | | | | |
| CDC | | v | | |
+---------+ | Write-Ahead Log | | |
| | | | |
| v | | |
| Micro-batch Buffer | ----> | Parquet |
| | | | files |
| v | | |
| Iceberg Commit | ----> | Metadata |
+---------------------+ +-------------+
- Streaming API: Records arrive via gRPC (Flight SQL), HTTP, or external connectors (such as Gnok CDC). The API validates schema compatibility and assigns WAL sequence numbers.
- Write-Ahead Log (WAL): Each record is durably written to the WAL before being acknowledged to the producer. This guarantees no data loss even if the server crashes before flushing to Iceberg.
- Micro-batch Buffer: Records accumulate in memory until a flush threshold is reached (time, row count, or byte size). The buffer then writes a Parquet file to object storage.
- Iceberg Commit: After the Parquet file is written, a new Iceberg snapshot is committed atomically. The commit references the new data file and becomes visible to all readers.
Data written to the WAL is immediately queryable within the producing session (read-your-writes), even before the Iceberg commit completes. Other sessions see data after the Iceberg commit.
Supported Sources
| Source | Status | Description |
|---|---|---|
| Direct Ingest | GA | Arrow RecordBatch ingestion via Flight SQL or HTTP |
| Gnok CDC | GA | MySQL CDC via Debezium/Kafka with MERGE upserts, schema evolution, and deduplication |
| COPY INTO | GA | Bulk loading from staged files (Parquet, CSV, JSON) with distributed execution |
Streaming Tables
Create a table with streaming enabled:
CREATE TABLE real_time_events (
event_id INT,
event_type STRING,
timestamp TIMESTAMP,
payload STRING
) STREAMING ENABLED;
-- Data is immediately queryable during ingestion
SELECT COUNT(*) FROM real_time_events WHERE event_type = 'click';
Write-Ahead Log (WAL)
The WAL provides durability for streaming data. When enabled, every record is persisted to the WAL before the producer receives an acknowledgment. If the server crashes, recovery replays uncommitted WAL entries to restore buffered data.
WAL Configuration
Durability and WAL storage are managed by Gnok. Customers do not configure service filesystem paths or switch the hosted service to an unsafe durability mode.
WAL Recovery
On startup, Gnok scans WAL segments for uncommitted entries (records written to WAL but not yet flushed to Iceberg). These entries are replayed into the micro-batch buffer and flushed normally. Recovery is automatic and requires no manual intervention.
WAL sequence numbers are used to detect and skip duplicate entries during recovery, ensuring exactly-once semantics even after crash and restart.
Micro-Batch Configuration
Service-managed thresholds determine when buffered records flush. Monitor ingestion progress and query-visible output instead of assuming an acknowledged source event is already in a committed table snapshot.
Backpressure
When producers write faster than the system can flush, backpressure prevents unbounded memory growth:
- Buffer threshold: When the in-memory buffer reaches the service flush threshold, a flush is triggered and new writes are queued.
- Queue depth limit: If the write queue exceeds a configurable depth (default: 4x the flush buffer size), the streaming API begins rejecting writes.
- Producer notification: gRPC producers receive flow-control signals via gRPC stream back-pressure. HTTP producers receive
429 Too Many Requestswith aRetry-Afterheader. TheRESOURCE_EXHAUSTEDgRPC status code indicates the server is overloaded. - Recovery: As flushes complete and buffer space is freed, the queue drains and new writes are accepted.
Backpressure is applied per-table. A slow-flushing table does not block ingestion into other tables.
Exactly-Once Semantics
Gnok guarantees exactly-once delivery from WAL write through Iceberg commit:
- WAL sequence numbers: Each record receives a monotonically increasing sequence number. If a producer retries a write (e.g., after a timeout), the WAL detects duplicates by sequence number and skips them.
- Idempotent Iceberg commits: Each flush operation is tagged with a unique commit ID derived from the WAL sequence range. If a commit is retried (e.g., after a transient S3 error), the Iceberg catalog's optimistic concurrency control detects the duplicate and rejects it without creating a duplicate snapshot.
- End-to-end guarantee: Producers that use WAL-backed streaming with retry can achieve exactly-once semantics without application-level deduplication.
Queryability
Data is queryable at two levels:
- Session-local (WAL): Within the producing session, records are readable as soon as they are written to the WAL. This provides read-your-writes consistency with sub-millisecond latency.
- Global (Iceberg): All sessions see data after the Iceberg snapshot commit. Latency depends on the flush interval (default 10 seconds).
-- In the producing session: sees WAL data immediately
INSERT INTO events VALUES (1, 'click', NOW(), '{}');
SELECT * FROM events WHERE event_id = 1; -- returns the row
-- In another session: sees data after the next Iceberg commit
SELECT * FROM events WHERE event_id = 1; -- visible after flush
CDC Integration
Gnok CDC replicates MySQL tables into Gnok Iceberg tables in real time via Debezium and Kafka.
Pipeline Architecture
MySQL (binlog) --> Debezium --> Kafka --> Gnok CDC --> Gnok
Setup Steps
CDC isn't self-service yet. Email support@gnok.io with your source database and the target catalog and schema; your administrator provides source access and grants on the target. Gnok runs the service-side ingestion components.
Schema Mapping
Gnok CDC automatically maps MySQL types to Gnok/Iceberg types:
| MySQL Type | Gnok Type |
|---|---|
INT, BIGINT | INT, BIGINT |
VARCHAR, TEXT | STRING |
DECIMAL(p,s) | DECIMAL(p,s) |
DATETIME, TIMESTAMP | TIMESTAMP |
DATE | DATE |
BOOLEAN, TINYINT(1) | BOOLEAN |
JSON | STRING |
New columns added to MySQL are automatically added to the Gnok table (schema evolution). Dropped columns are retained in the Iceberg schema as nullable.
Deduplication
Gnok CDC deduplicates within each micro-batch using the primary key. For entity tables (tables with a primary key), it uses MERGE statements to upsert. For append-only audit tables (no primary key), it uses INSERT.
Monitoring
View active streams and their metrics:
SHOW STREAMS;
Returns:
| Column | Description |
|---|---|
stream_name | Identifier for the stream (table name or connector name) |
table_name | Target Iceberg table |
status | ACTIVE, PAUSED, or ERROR |
rows_ingested | Total rows ingested since stream creation |
bytes_written | Total bytes written to Parquet |
lag_ms | Time since the last record was received (indicates producer health) |
last_flush_time | Timestamp of the most recent Iceberg commit |
buffer_rows | Rows currently in the micro-batch buffer |
buffer_bytes | Bytes currently in the micro-batch buffer |
wal_sequence | Current WAL sequence number |
Health Checks
Monitor these indicators:
lag_msincreasing: The producer may be down or the network path is degraded.buffer_rowsconsistently high: Ingestion may be outpacing flushes. Flush thresholds are service-managed; contact support@gnok.io if it persists.status = ERROR: Common causes include schema mismatches and object storage errors. Email support@gnok.io with the stream name and the time the error started.
Further Reading
- COPY INTO -- bulk file ingestion
- Table Maintenance -- compacting streamed data