Skip to content

Replace the logical codec's object registry with durable metadata #1724

Description

@timsaucer

examples/datafusion-ffi-example/src/logical_extension_codec.rs parks live table providers in a process-global HashMap and encodes an integer token into it. Encoding inserts, decoding removes, so the same bytes cannot be decoded twice, one plan cannot fan out to several readers, and a plan that never reaches a decoder keeps its provider alive for the life of the process. extension-guide/codecs.md tells authors not to do this.

Unlike the physical codec in the same crate (see the quarantine sub-issue), this one is fixable: try_encode_table_provider at line 148 claims node.downcast_ref::<MemTable>(), which is narrow, and a MemTable is fully describable by its schema and batches.

Pattern to copy: examples/distributed/storage-library/src/codec.rs — same Arrow IPC technique, same error convention (internal_datafusion_err! on encode, since this process holds the object; exec_datafusion_err! on decode, since those are foreign bytes).

Verified prerequisites: MemTable.batches is pub (datafusion-catalog/src/memory/table.rs:69), typed Vec<PartitionData> where PartitionData = Arc<tokio::sync::RwLock<Vec<RecordBatch>>>. MemTable::try_new rejects zero partitions (table.rs:84). arrow is already a dependency with IPC available, so no Cargo.toml change.

Proposed wire format, keeping the per-instance prefix the dispatch tests rely on:

<provider_prefix> | b"MEMTBL1" | u32 LE n_partitions | { u32 LE ipc_len | ipc stream }*

One stream per partition, because MemTable partition boundaries become output partitions. A stream carries its schema even when empty, so an empty partition round-trips.

Two traps worth writing down before someone hits them:

  • Use try_read() on each partition lock, not blocking_read(). The FFI codec runs with a tokio runtime handle installed, and blocking_read panics in that context.
  • On decode, build MemTable::try_new from the IPC schema, not the schema: SchemaRef argument. try_new validates schema.contains(&batch.schema()), so metadata drift would surface as a spurious mismatch. This is the opposite choice from storage-library/src/codec.rs:376-380, which must honour the plan's schema because it re-reads files from disk; here the batches are the payload. Worth a comment noting the contrast, since the two codecs otherwise look alike.

Done when: the registry, token_id(), and the HashMap/Mutex/OnceLock/AtomicU64 imports are gone; the struct field token is renamed provider_prefix to match the Python kwarg that already uses that name; and grep -in token over the file returns nothing.

Tests: of 19 tests in python/tests/_test_logical_extension_codec.py, one changes. test_installing_a_codec_cannot_hijack_an_earlier_codecs_objects asserts len(before) == len(after) with a comment about tokens being minted per encode; that comment becomes false and the assertion becomes weaker than reality, so it should become assert before == after. Add one test for the property the guide claims and nothing currently covers: encode once, decode twice on one session, assert both produce the same rows. All 47 tests in the planner crate should be unaffected — every assertion there is on call counters, never on payload shape.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or requestrustPull requests that update Rust code

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions