Skip to content

take_blobs fails on the Ray runner: LanceBlobFile is not picklable #65

Description

@FANNG1

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:

  1. 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.
  2. 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).

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

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions