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]

Reply via email to