timsaucer opened a new issue, #1724:
URL: https://github.com/apache/datafusion-python/issues/1724

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


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to