Warning
This is not authoritative documentation. It describes a plan of work that is not yet implemented and will change as it lands. Once the work is complete this document should be removed; the finished system is described by the Iceberg storage engine spec, the plugin’s user documentation, and the repos themselves.
Task 4: Read path — planner, executor, cursor, pinning¶
Repo: https://opendev.org/drizzle/drizzle. Depends on task 3. This is the core read deliverable: after it, use case 1 (SQL over pipeline output) works, full-scan. Three commits.
Design constraints from the design doc, restated as binding:
Snapshot resolution is an explicit planning step producing a value.
Planner output must be constructible from a received task list as well as from local planning (the future distribution seam) — enforce by giving the Planner a constructor/factory taking
vector<FileScanTask>directly, used by tests today.Pinning is per transaction (autocommit → per statement). Never per cursor.
Commit 1: planner and executor (no cursor yet; unit-tested standalone)¶
planner.{h,cc}:(table handle, snapshot-id, projected iceberg field-ids)→ immutable orderedvector<FileScanTask>via iceberg-cpp scan planning (V2 deletes handled by the library). Plus the from-existing-tasks factory.executor.{h,cc}: one FileScanTask → row stream. Internally: Arrow record-batch reader with column projection by field-id; exposesnext_row()filling a caller-provided callback/row-view, plusseek(row_ordinal)used byrnd_pos(implement as re-open + skip if the reader lacks random access — correctness first; note the cost in a comment).Value conversion Arrow → Drizzle
Field::storelives here, the inverse of task 3’s schema map and colocated with it conceptually: cover every mapped type, with explicit UTC handling for timestamptz→EPOCH and exactness tests for decimal.
Commit 2: the cursor, wired¶
cursor.{h,cc} (IcebergCursor : Cursor), modeled on
plugin/function_engine/cursor.{h,cc} for shape:
doOpen: resolve the per-table share (schema map, field-id table).doStartTableScan(bool): obtain(snapshot-id, plan)through the session pin (below); position at task 0.rnd_next(buf): drain executors task by task; pack via the read map — honorTable::read_set(table.h:497) so unread columns are neither fetched from Parquet nor stored. AddHTON_PARTIAL_COLUMN_READto the engine flags in this commit.position(record): encode(task ordinal: u32, row ordinal: u64)intoref; setref_lengthaccordingly at open.rnd_pos(buf, pos): decode; executor seek within the pinned plan.info(flag):stats.recordsfrom the pinned snapshot’s summarytotal-records(estimate; the engine deliberately does not setHTON_STATS_RECORDS_IS_EXACT),stats.mean_rec_lengthfrom the proto.close()/doEndTableScan: release executors; the plan is owned by the session pin entry, not the cursor (two cursors on one table share one plan).No index methods implemented;
doStartIndexScannever reachable (no keys are ever declared in the synthesized proto).
Session pinning: session_state.{h,cc}¶
Per-session engine slot via
Session::getEngineData(const plugin::MonitoredInTransaction*)
(session.h:345; the innobase trx_t* pattern,
ha_innodb.cc:1030). Holds
map<table-uuid, PinEntry{snapshot-id, shared plan cache}>. First
touch of a table in a transaction resolves current-snapshot and pins;
subsequent touches (including second cursors in the same statement — the
self-join case) reuse the pin. Clearing: engine doEndStatement
(storage_engine.h:150-155) when no transaction is open (autocommit),
and transaction end otherwise — in this task the engine is not yet
transactional, so wire statement-end clearing now and leave a
// task 5 moves clearing to doCommit/doRollback when a transaction is open
marker that task 5 must resolve (no other tombstones).
Commit 3: test suites¶
Full-scan SELECT over each fixture table; checksums match pyiceberg-read values recorded in the fixtures (interop, not self-consistency).
The position-deletes fixture: deleted rows absent.
ORDER BY on a non-trivial table (exercises filesort →
position/rnd_pos).Self-join on a fixture table while a background pyiceberg append commits mid-statement (the suite’s seeding hook can do this): both sides join the same snapshot. Also: two statements in one autocommit session straddling an external commit see old-then-new (statement-level pinning observable).
Projection:
SELECT one_col FROM wide_fixture— assert via the callgrind CI job or a byte-counter that unprojected columns are not read (if neither is practical, a unit test on the executor’s projection set is the floor).
Verification¶
Build green in both switch states (
--with-icebergon/off) per commit; suites green; job voting.Valgrind server start/scan/stop: no leaks attributable to the plugin; massif comparison sane on the wide-table scan (row-at-a-time packing must not accumulate whole files in memory — batch-at-a-time residency only).