JingsongLi commented on code in PR #9130:
URL: https://github.com/apache/paimon/pull/9130#discussion_r4083014133
##########
paimon-python/pypaimon/ray/ray_paimon.py:
##########
@@ -303,6 +387,102 @@ def _map_blob_batch(
return result
+def _map_blob_affinity_block(
+ batch, file_io, blob_cols, all_blob_cols, parallelism, fn, fn_kwargs,
+ fn_batch_size, prefetch_bytes, affinity_cols,
+ map_blob_cols=(), array_blob_cols=()):
+ from pypaimon.multimodal.blob_read import fetch_blob_bodies
+
+ if batch.num_rows == 0:
+ return
+
+ scalar_cols = _blob_scalar_columns(
+ batch, blob_cols, all_blob_cols, affinity_cols)
+
+ for start, end in _blob_prefetch_windows(
+ batch, blob_cols, fn_batch_size, prefetch_bytes):
+ window = batch.slice(start, end - start)
+ bodies = fetch_blob_bodies(
+ file_io,
+ window.select(blob_cols).to_pydict(),
+ blob_cols,
+ parallelism,
+ )
Review Comment:
[P2] Release the completed payload window before fetching the next one
On the second and subsequent iterations, the previous `bodies` remains
referenced while the right-hand side of this assignment fetches the entire next
window. `fn_bodies` also retains the previous window's last inference batch.
Consequently, even a UDF that only returns small scalar results keeps two
payload windows alive during the next read, adding roughly 64 MiB per active
worker for full default-sized windows. With real `LocalFileIO`, 1 MiB blobs,
`batch_size=1`, and a 2 MiB window, I measured about 6 MiB peak allocation;
releasing the completed window reduced it to about 4 MiB with the same
coalesced-read buffers. Please release both `bodies` and `fn_bodies` after
consuming each window, before starting the next fetch.
##########
paimon-python/pypaimon/ray/ray_paimon.py:
##########
@@ -225,19 +240,88 @@ def map_with_blobs(
or (map_blob_columns | array_blob_columns) - all_blob):
raise ValueError("Nested BLOB columns must be disjoint subsets of
all_blob_columns.")
+ mapper = _map_blob_batch
+ affinity_cols = []
+ if blob_uri_affinity:
+ if (map_blob_columns | array_blob_columns).intersection(blob_cols):
+ raise ValueError("blob_uri_affinity supports scalar BLOB columns
only")
+ dataset, affinity_cols = _cluster_by_blob_uri(dataset, blob_cols)
+ mapper = _map_blob_affinity_block
+ kwargs["batch_size"] = None
Review Comment:
[P2] Preserve a valid Ray batch size when requesting GPUs
Setting the outer `batch_size` to `None` makes affinity mode fail for GPU
inference. For example, `map_with_blobs(..., batch_size=32,
blob_uri_affinity=True, num_gpus=1)` immediately raises `ValueError: You must
provide batch_size to map_batches when requesting GPUs`, before the UDF runs.
Passing `ray_remote_args={"num_gpus": 1}` has the same result; both calls are
accepted with affinity disabled. This is reproducible with Ray 2.56.1's actual
`Dataset.map_batches` validation and also affects the declared minimum Ray
version. Please use an outer batching strategy compatible with Ray's GPU
validation while retaining the internal inference batches, and add coverage for
GPU resource arguments.
--
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]