# @ray.remote function seemingly copying data from plasma store

**URL:** <https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297>\
**Category:** Ray Core\
**Created:** [March 17, 2021, 4:12pm UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297 "2021-03-17T16:12:46Z")\
**Posts on this page:** 11\
**Page:** 1

<div class="post-metadata">

**Author:** ![Robert\_Speare](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/robert_speare/32/641_2.png) [@Robert\_Speare](https://discuss.ray.io/u/Robert_Speare)\
**Post date:** [March 17, 2021, 4:12pm UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/1 "2021-03-17T16:12:46Z")

</div>

Hello,

I am currently trying to reduce the memory of a program that takes one large numpy array and one large dictionary of pandas dataframes, and iterates over their values using different model parameters to create some score/result. I would like to share these two objects across processes without copying. After configuring my model function to be remote, and passing references to the large numpy/pandas objects and running things with the [mprof](https://github.com/pythonprofilers/memory_profiler/blob/master/mprof.py) memory profiler, I am finding that memory overhead for the program is the same as running with ProcessPool, where these objects are piped/copied to each process. What am I missing to get to zero-copy? (Possibly related thread: [How to share memory with non-numpy object?](https://discuss.ray.io/t/how-to-share-memory-with-non-numpy-object/1295))

The pattern looks like this:

```python
@ray.remote
def run_model(array, dict_of_pandas_df, params):
     ... do some work ...
     return result

numpy_array_ref = ray.put(my_array)
dict_of_pandas_df_ref = ray.put(dict_of_pandas_df)
list_of_model_parameters = [{}, {}...]
result_refs = [run_model.remote(numpy_array_ref, dict_of_pandas_df_ref, x) for x in list_of_model_parameters]

results = ray.get(result_refs)

```

---

<div class="post-metadata">

**Author:** ![sangcho](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/sangcho/32/425_2.png) [@sangcho](https://discuss.ray.io/u/sangcho)\
**Post date:** [March 17, 2021, 5:41pm UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/2 "2021-03-17T17:41:12Z")

</div>

When you do zero-copy to your process, note that the process still says it uses the memory (from the shared memory). It means your memory is double counted. For example, if you have 2 processes A and B, each of which uses 100MB of shared memory, both of them will say they use 100MB of memory.

cc @suquark do you know any good way to verify the zero copy read?

---

<div class="post-metadata">

**Author:** ![Robert\_Speare](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/robert_speare/32/641_2.png) [@Robert\_Speare](https://discuss.ray.io/u/Robert_Speare)\
**Post date:** [March 17, 2021, 7:03pm UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/3 "2021-03-17T19:03:50Z")

</div>

Thanks a lot @sangcho . If it’s useful, I am profiling with the following mprofile command:

```bash
mprof run --include-children python my_script.py

```

Even when looking at the running processes in `top`/`htop`, I believe that the RES portion of memory is significantly larger than SHR – both when running multiprocessing as well as with `@ray.remote` (will follow up to verify on this)

---

<div class="post-metadata">

**Author:** ![sangcho](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/sangcho/32/425_2.png) [@sangcho](https://discuss.ray.io/u/sangcho)\
**Post date:** [March 17, 2021, 11:20pm UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/4 "2021-03-17T23:20:42Z")

</div>

Hmm I see. What’s the dtype of your numpy. Are they float or integer or strings?

---

<div class="post-metadata">

**Author:** ![Robert\_Speare](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/robert_speare/32/641_2.png) [@Robert\_Speare](https://discuss.ray.io/u/Robert_Speare)\
**Post date:** [March 18, 2021, 12:12am UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/5 "2021-03-18T00:12:37Z")

</div>

I am passing in a numpy array with dtype float. The pandas dataframes – values of the dict – have mixed types: float/pandas.Timestamp.

---

<div class="post-metadata">

**Author:** ![sangcho](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/sangcho/32/425_2.png) [@sangcho](https://discuss.ray.io/u/sangcho)\
**Post date:** [March 18, 2021, 12:33am UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/6 "2021-03-18T00:33:00Z")

</div>

Hmm that’s pretty weird. Pandas dataframe uses the numpy array under the hood afaik, and float type numpy should be zero-copy read. So each process only should copy parts that are not numpy array which means your SHR should be larger than RES)!

Can you actually try sth like this?

```auto
@ray.remote
def your_func(array_ref_list):
    # Measure the overhead of ray.get
    s = time.perf_counter()
    ray.get(array_ref_list)
    print(time.perf_counter() - s)

numpy_array_ref = ray.put(my_array)
# NOTE: Pass the list of object ref so that it won't be automatically obtained.
[your_func([numpy_array_ref]) for _ in range(10)]

```

---

<div class="post-metadata">

**Author:** ![sangcho](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/sangcho/32/425_2.png) [@sangcho](https://discuss.ray.io/u/sangcho)\
**Post date:** [March 18, 2021, 12:33am UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/7 "2021-03-18T00:33:33Z")

</div>

If the overhead is as big as copying large array into the process, that means zero copy read wasn’t working as expected.

---

<div class="post-metadata">

**Author:** ![sangcho](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/sangcho/32/425_2.png) [@sangcho](https://discuss.ray.io/u/sangcho)\
**Post date:** [March 18, 2021, 12:38am UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/8 "2021-03-18T00:38:17Z")

</div>

Actually there’s also a possibility `pandas.Timestamp.` is not zero-copyable.

---

<div class="post-metadata">

**Author:** ![Robert\_Speare](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/robert_speare/32/641_2.png) [@Robert\_Speare](https://discuss.ray.io/u/Robert_Speare)\
**Post date:** [March 18, 2021, 2:07am UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/9 "2021-03-18T02:07:47Z")

</div>

Awesome @sangcho. I’ll try what you’ve put forward – might take 24 hours or so to get back on this thread. Thanks!

---

<div class="post-metadata">

**Author:** ![Robert\_Speare](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/robert_speare/32/641_2.png) [@Robert\_Speare](https://discuss.ray.io/u/Robert_Speare)\
**Post date:** [March 24, 2021, 1:10am UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/10 "2021-03-24T01:10:16Z")

</div>

Hey @sangcho. After looking into things, the latency of a `ray.get` with a function is not very large – only a few seconds – for an object of ~50 GB. Need to do some more research, but hoping to post back on this thread with a fully reproducible example if this rears its head again – thank you so much for your help!

---

<div class="post-metadata">

**Author:** ![sangcho](https://sea2.discourse-cdn.com/flex020/user_avatar/discuss.ray.io/sangcho/32/425_2.png) [@sangcho](https://discuss.ray.io/u/sangcho)\
**Post date:** [March 27, 2021, 10:31pm UTC](https://discuss.ray.io/t/ray-remote-function-seemingly-copying-data-from-plasma-store/1297/11 "2021-03-27T22:31:52Z")

</div>

Sounds good! I still have suspicion that the issue is you have Timestamp dtype btw (afaik, we only support zero-copy read for integer & float).
