Skip to main content

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 |
+---------------------+ +-------------+
  1. 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.
  2. 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.
  3. 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.
  4. 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​

SourceStatusDescription
Direct IngestGAArrow RecordBatch ingestion via Flight SQL or HTTP
Gnok CDCGAMySQL CDC via Debezium/Kafka with MERGE upserts, schema evolution, and deduplication
COPY INTOGABulk 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:

  1. Buffer threshold: When the in-memory buffer reaches the service flush threshold, a flush is triggered and new writes are queued.
  2. Queue depth limit: If the write queue exceeds a configurable depth (default: 4x the flush buffer size), the streaming API begins rejecting writes.
  3. Producer notification: gRPC producers receive flow-control signals via gRPC stream back-pressure. HTTP producers receive 429 Too Many Requests with a Retry-After header. The RESOURCE_EXHAUSTED gRPC status code indicates the server is overloaded.
  4. 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 TypeGnok Type
INT, BIGINTINT, BIGINT
VARCHAR, TEXTSTRING
DECIMAL(p,s)DECIMAL(p,s)
DATETIME, TIMESTAMPTIMESTAMP
DATEDATE
BOOLEAN, TINYINT(1)BOOLEAN
JSONSTRING

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:

ColumnDescription
stream_nameIdentifier for the stream (table name or connector name)
table_nameTarget Iceberg table
statusACTIVE, PAUSED, or ERROR
rows_ingestedTotal rows ingested since stream creation
bytes_writtenTotal bytes written to Parquet
lag_msTime since the last record was received (indicates producer health)
last_flush_timeTimestamp of the most recent Iceberg commit
buffer_rowsRows currently in the micro-batch buffer
buffer_bytesBytes currently in the micro-batch buffer
wal_sequenceCurrent WAL sequence number

Health Checks​

Monitor these indicators:

  • lag_ms increasing: The producer may be down or the network path is degraded.
  • buffer_rows consistently 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​