Problem
daft_lance.take_blobs fails on the Ray runner. Any query that materializes its output fails with:
TypeError: cannot pickle 'builtins.LanceBlobFile' object
The native runner works, so the existing tests pass locally and in CI. CI never runs under Ray, since ray is not part of the locked dev environment.
The behavior is the same on pylance 8.0.0 and 11.0.0, so it is not a regression from a Lance upgrade.
Reproduction
import lance, pyarrow as pa, daft
from daft_lance import take_blobs
daft.set_runner_ray()
table = pa.table({"id": [1, 2, 3], "blob": lance.blob_array([b"a", b"b", b"c"])})
ds = lance.write_dataset(table, "/tmp/blobs.lance", data_storage_version="2.2")
df = daft.read_lance(ds.uri, default_scan_options={"with_row_id": True})
take_blobs(df, ds, "blob").to_pydict() # TypeError: cannot pickle 'builtins.LanceBlobFile' object
Running tests/io/lancedb/test_lance_blob_v2_read.py with DAFT_RUNNER=ray fails every test that collects take_blobs output (test_take_blobs_kinds, test_take_blobs_single_row, test_take_blobs_all_rows, test_take_blobs_non_contiguous_rows).
Root cause
The take_blobs UDF (daft_lance/_blob.py) returns the lance.BlobFile objects from LanceDataset.take_blobs as a DataType.python() column. On Ray, a UDF's output partition is serialized with cloudpickle to move between actors and back to the driver. BlobFile wraps a Rust LanceBlobFile handle that is not picklable, so the task fails when its results are serialized.
Even if the handle were picklable, a BlobFile is a one-shot stream bound to an open dataset in the worker process (see #42), so it can't be carried across process boundaries anyway.
Expected behavior
take_blobs should give the same results on the native and Ray runners, including rows whose blob is null. Since pylance 10, LanceDataset.take_blobs returns None for null blobs and keeps results aligned with the requested row ids.
Proposed fix
Pick one and apply it consistently:
- Materialize bytes in the UDF. Call
read() on each BlobFile inside the worker and return a binary column (null for null blobs). This is simple and picklable, but it gives up lazy/ranged reads for large blobs.
- Return a picklable lazy handle. Return a small object that holds the dataset URI, version, storage options, column and row id, and reopens the dataset to produce a
BlobFile on first access. This keeps lazy reads and works across processes, but it is a larger API change.
Either way, the UDF should reopen the dataset from a serializable context (like the DatasetOpenContext used by the maintenance UDFs) instead of closing over a live LanceDataset.
Regression test
Run the take_blobs tests on both runners. Either parametrize them over native and Ray, or add a Ray CI job with ray installed. Assert that:
- every storage kind (inline, packed, dedicated, external full, external slice) reads back the expected bytes;
- non-contiguous row selections read back the expected bytes;
- null blobs come back as null and stay aligned with their rows.
Related: #30 (Blob V2 tracking), #42 (one-shot BlobFile semantics).
Problem
daft_lance.take_blobsfails on the Ray runner. Any query that materializes its output fails with:The native runner works, so the existing tests pass locally and in CI. CI never runs under Ray, since
rayis not part of the locked dev environment.The behavior is the same on pylance 8.0.0 and 11.0.0, so it is not a regression from a Lance upgrade.
Reproduction
Running
tests/io/lancedb/test_lance_blob_v2_read.pywithDAFT_RUNNER=rayfails every test that collectstake_blobsoutput (test_take_blobs_kinds,test_take_blobs_single_row,test_take_blobs_all_rows,test_take_blobs_non_contiguous_rows).Root cause
The
take_blobsUDF (daft_lance/_blob.py) returns thelance.BlobFileobjects fromLanceDataset.take_blobsas aDataType.python()column. On Ray, a UDF's output partition is serialized with cloudpickle to move between actors and back to the driver.BlobFilewraps a RustLanceBlobFilehandle that is not picklable, so the task fails when its results are serialized.Even if the handle were picklable, a
BlobFileis a one-shot stream bound to an open dataset in the worker process (see #42), so it can't be carried across process boundaries anyway.Expected behavior
take_blobsshould give the same results on the native and Ray runners, including rows whose blob is null. Since pylance 10,LanceDataset.take_blobsreturnsNonefor null blobs and keeps results aligned with the requested row ids.Proposed fix
Pick one and apply it consistently:
read()on eachBlobFileinside the worker and return abinarycolumn (null for null blobs). This is simple and picklable, but it gives up lazy/ranged reads for large blobs.BlobFileon first access. This keeps lazy reads and works across processes, but it is a larger API change.Either way, the UDF should reopen the dataset from a serializable context (like the
DatasetOpenContextused by the maintenance UDFs) instead of closing over a liveLanceDataset.Regression test
Run the
take_blobstests on both runners. Either parametrize them over native and Ray, or add a Ray CI job withrayinstalled. Assert that:Related: #30 (Blob V2 tracking), #42 (one-shot
BlobFilesemantics).