mars-project / mars-project/mars

[BUG]failed call mars.tensor.sum in mars.dataframe.apply

Open
#2,219 1 comment 0 reactions 0 assignees View on GitHub
mod: dataframe type: bug
Dominant language
Python
Stars
2.7k
Forks
325
PR merge metrics
No merged PRs in 30d

Description

```python3
def mfunc(x):
return mt.sum(x)

mdf = md.DataFrame(mt.random.randn(5,6))
res = mdf.apply(func=mfunc,axis=1).execute()
print(res)
```

```shell

2021-07-15 21:55:09,494 - mars.services.scheduling.worker.execution - ERROR - Failed to run subtask wYZbazCTIG4yPjL1SOmdWTdq on band numa-0
Traceback (most recent call last):
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/scheduling/worker/execution.py", line 223, in internal_run_subtask
yield asyncio.shield(run_aiotask)
File "mars/oscar/core.pyx", line 264, in mars.oscar.core._Actor._run_actor_async_generator
File "mars/oscar/core.pyx", line 206, in _handle_actor_result
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/subtask/api.py", line 52, in run_subtask_in_slot
return await ref.run_subtask(subtask)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 151, in send
return self._process_result_message(result)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 59, in _process_result_message
raise message.error.with_traceback(message.traceback)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/pool.py", line 482, in send
result = await future
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/api.py", line 111, in __on_receive__
return await super().__on_receive__(message)
File "mars/oscar/core.pyx", line 321, in __on_receive__
File "mars/oscar/core.pyx", line 315, in mars.oscar.core._Actor.__on_receive__
File "mars/oscar/core.pyx", line 237, in _handle_actor_result
File "mars/oscar/core.pyx", line 258, in mars.oscar.core._Actor._run_actor_async_generator
File "mars/oscar/core.pyx", line 262, in mars.oscar.core._Actor._run_actor_async_generator
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/subtask/worker/runner.py", line 85, in run_subtask
result = yield self._running_processor.run(subtask)
File "mars/oscar/core.pyx", line 264, in mars.oscar.core._Actor._run_actor_async_generator
File "mars/oscar/core.pyx", line 206, in _handle_actor_result
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 151, in send
return self._process_result_message(result)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 59, in _process_result_message
raise message.error.with_traceback(message.traceback)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/pool.py", line 482, in send
result = await future
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/api.py", line 111, in __on_receive__
return await super().__on_receive__(message)
File "mars/oscar/core.pyx", line 321, in __on_receive__
File "mars/oscar/core.pyx", line 315, in mars.oscar.core._Actor.__on_receive__
File "mars/oscar/core.pyx", line 237, in _handle_actor_result
File "mars/oscar/core.pyx", line 258, in mars.oscar.core._Actor._run_actor_async_generator
File "mars/oscar/core.pyx", line 262, in mars.oscar.core._Actor._run_actor_async_generator
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/subtask/worker/processor.py", line 424, in run
result = yield self._running_aio_task
File "mars/oscar/core.pyx", line 264, in mars.oscar.core._Actor._run_actor_async_generator
File "mars/oscar/core.pyx", line 206, in _handle_actor_result
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/subtask/worker/processor.py", line 263, in run
store_infos = await put_infos
File "/root/miniconda3/lib/python3.7/site-packages/mars/utils.py", line 1152, in _async_batch
return await self.batch_func(args_list, kwargs_list)
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/storage/api/oscar.py", line 109, in batch_put
return await self._storage_handler_ref.put.batch(*puts)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 151, in send
return self._process_result_message(result)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 59, in _process_result_message
raise message.error.with_traceback(message.traceback)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/pool.py", line 482, in send
result = await future
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/api.py", line 111, in __on_receive__
return await super().__on_receive__(message)
File "mars/oscar/core.pyx", line 321, in __on_receive__
File "mars/oscar/core.pyx", line 305, in mars.oscar.core._Actor.__on_receive__
File "mars/oscar/core.pyx", line 311, in mars.oscar.core._Actor.__on_receive__
File "/root/miniconda3/lib/python3.7/site-packages/mars/utils.py", line 1152, in _async_batch
return await self.batch_func(args_list, kwargs_list)
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/storage/core.py", line 336, in batch_put
object_info = await self._clients[level].put(obj)
File "/root/miniconda3/lib/python3.7/site-packages/mars/storage/shared_memory.py", line 150, in put
buffers = await serializer.run()
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/aio.py", line 59, in run
return self._get_buffers()
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/aio.py", line 34, in _get_buffers
headers, buffers = serialize(self._obj)
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py", line 339, in serialize
result = serializer.serialize(obj, context)
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py", line 75, in wrapped
return func(self, obj, context)
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py", line 152, in serialize
return {}, pickle_buffers(obj)
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py", line 92, in pickle_buffers
protocol=BUFFER_PICKLE_PROTOCOL,
File "/root/miniconda3/lib/python3.7/site-packages/cloudpickle/cloudpickle_fast.py", line 73, in dumps
cp.dump(obj)
File "/root/miniconda3/lib/python3.7/site-packages/cloudpickle/cloudpickle_fast.py", line 563, in dump
return Pickler.dump(self, obj)
TypeError: can't pickle weakref objects
2021-07-15 21:55:09,519 - asyncio - ERROR - Exception in callback _execute.._attach_session() at /root/miniconda3/lib/python3.7/site-packages/mars/core/session.py:611
handle: ._attach_session() at /root/miniconda3/lib/python3.7/site-packages/mars/core/session.py:611>
Traceback (most recent call last):
File "/root/miniconda3/lib/python3.7/asyncio/events.py", line 88, in _run
self._context.run(self._callback, *self._args)
File "/root/miniconda3/lib/python3.7/site-packages/mars/core/session.py", line 612, in _attach_session
fut.result()
File "/root/miniconda3/lib/python3.7/site-packages/mars/deploy/oscar/session.py", line 160, in _run_in_background
raise task_result.error.with_traceback(task_result.traceback)
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/scheduling/worker/execution.py", line 223, in internal_run_subtask
yield asyncio.shield(run_aiotask)
File "mars/oscar/core.pyx", line 264, in mars.oscar.core._Actor._run_actor_async_generator
File "mars/oscar/core.pyx", line 206, in _handle_actor_result
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/subtask/api.py", line 52, in run_subtask_in_slot
return await ref.run_subtask(subtask)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 151, in send
return self._process_result_message(result)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 59, in _process_result_message
raise message.error.with_traceback(message.traceback)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/pool.py", line 482, in send
result = await future
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/api.py", line 111, in __on_receive__
return await super().__on_receive__(message)
File "mars/oscar/core.pyx", line 321, in __on_receive__
File "mars/oscar/core.pyx", line 315, in mars.oscar.core._Actor.__on_receive__
File "mars/oscar/core.pyx", line 237, in _handle_actor_result
File "mars/oscar/core.pyx", line 258, in mars.oscar.core._Actor._run_actor_async_generator
File "mars/oscar/core.pyx", line 262, in mars.oscar.core._Actor._run_actor_async_generator
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/subtask/worker/runner.py", line 85, in run_subtask
result = yield self._running_processor.run(subtask)
File "mars/oscar/core.pyx", line 264, in mars.oscar.core._Actor._run_actor_async_generator
File "mars/oscar/core.pyx", line 206, in _handle_actor_result
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 151, in send
return self._process_result_message(result)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 59, in _process_result_message
raise message.error.with_traceback(message.traceback)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/pool.py", line 482, in send
result = await future
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/api.py", line 111, in __on_receive__
return await super().__on_receive__(message)
File "mars/oscar/core.pyx", line 321, in __on_receive__
File "mars/oscar/core.pyx", line 315, in mars.oscar.core._Actor.__on_receive__
File "mars/oscar/core.pyx", line 237, in _handle_actor_result
File "mars/oscar/core.pyx", line 258, in mars.oscar.core._Actor._run_actor_async_generator
File "mars/oscar/core.pyx", line 262, in mars.oscar.core._Actor._run_actor_async_generator
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/subtask/worker/processor.py", line 424, in run
result = yield self._running_aio_task
File "mars/oscar/core.pyx", line 264, in mars.oscar.core._Actor._run_actor_async_generator
File "mars/oscar/core.pyx", line 206, in _handle_actor_result
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/subtask/worker/processor.py", line 263, in run
store_infos = await put_infos
File "/root/miniconda3/lib/python3.7/site-packages/mars/utils.py", line 1152, in _async_batch
return await self.batch_func(args_list, kwargs_list)
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/storage/api/oscar.py", line 109, in batch_put
return await self._storage_handler_ref.put.batch(*puts)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 151, in send
return self._process_result_message(result)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py", line 59, in _process_result_message
raise message.error.with_traceback(message.traceback)
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/pool.py", line 482, in send
result = await future
File "/root/miniconda3/lib/python3.7/site-packages/mars/oscar/api.py", line 111, in __on_receive__
return await super().__on_receive__(message)
File "mars/oscar/core.pyx", line 321, in __on_receive__
File "mars/oscar/core.pyx", line 305, in mars.oscar.core._Actor.__on_receive__
File "mars/oscar/core.pyx", line 311, in mars.oscar.core._Actor.__on_receive__
File "/root/miniconda3/lib/python3.7/site-packages/mars/utils.py", line 1152, in _async_batch
return await self.batch_func(args_list, kwargs_list)
File "/root/miniconda3/lib/python3.7/site-packages/mars/services/storage/core.py", line 336, in batch_put
object_info = await self._clients[level].put(obj)
File "/root/miniconda3/lib/python3.7/site-packages/mars/storage/shared_memory.py", line 150, in put
buffers = await serializer.run()
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/aio.py", line 59, in run
return self._get_buffers()
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/aio.py", line 34, in _get_buffers
headers, buffers = serialize(self._obj)
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py", line 339, in serialize
result = serializer.serialize(obj, context)
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py", line 75, in wrapped
return func(self, obj, context)
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py", line 152, in serialize
return {}, pickle_buffers(obj)
File "/root/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py", line 92, in pickle_buffers
protocol=BUFFER_PICKLE_PROTOCOL,
File "/root/miniconda3/lib/python3.7/site-packages/cloudpickle/cloudpickle_fast.py", line 73, in dumps
cp.dump(obj)
File "/root/miniconda3/lib/python3.7/site-packages/cloudpickle/cloudpickle_fast.py", line 563, in dump
return Pickler.dump(self, obj)
TypeError: can't pickle weakref objects
---------------------------------------------------------------------------
TypeError Traceback (most recent call last)
in
3
4 mdf = md.DataFrame(mt.random.randn(5,6))
----> 5 res = mdf.apply(func=mfunc,axis=1).execute()
6 print(res)

~/miniconda3/lib/python3.7/site-packages/mars/core/entity/tileables.py in execute(self, session, **kw)
430
431 def execute(self, session=None, **kw):
--> 432 result = self.data.execute(session=session, **kw)
433 if isinstance(result, TILEABLE_TYPE):
434 return self

~/miniconda3/lib/python3.7/site-packages/mars/core/entity/executable.py in execute(self, session, **kw)
78
79 session = _get_session(self, session)
---> 80 return execute(self, session=session, **kw)
81
82 def _check_session(self,

~/miniconda3/lib/python3.7/site-packages/mars/core/session.py in execute(tileable, session, wait, backend, new_session_kwargs, show_progress, progress_update_interval, *tileables, **kwargs)
679 show_progress=show_progress,
680 progress_update_interval=progress_update_interval,
--> 681 **kwargs)
682
683

~/miniconda3/lib/python3.7/site-packages/mars/core/session.py in _inner(*args, **kwargs)
742 fut = _pool.submit(run_in_thread)
743 try:
--> 744 result, default_session_in_thread = fut.result()
745 except KeyboardInterrupt: # pragma: no cover
746 logger.warning('Cancelling running task')

~/miniconda3/lib/python3.7/concurrent/futures/_base.py in result(self, timeout)
433 raise CancelledError()
434 elif self._state == FINISHED:
--> 435 return self.__get_result()
436 else:
437 raise TimeoutError()

~/miniconda3/lib/python3.7/concurrent/futures/_base.py in __get_result(self)
382 def __get_result(self):
383 if self._exception:
--> 384 raise self._exception
385 else:
386 return self._result

~/miniconda3/lib/python3.7/concurrent/futures/thread.py in run(self)
55
56 try:
---> 57 result = self.fn(*self.args, **self.kwargs)
58 except BaseException as exc:
59 self.future.set_exception(exc)

~/miniconda3/lib/python3.7/site-packages/mars/core/session.py in run_in_thread()
738 # set default session in this thread
739 _sync_default_session(default_session)
--> 740 return func(*args, **kwargs), get_default_session()
741
742 fut = _pool.submit(run_in_thread)

~/miniconda3/lib/python3.7/site-packages/mars/core/session.py in execute(self, tileable, show_progress, *tileables, **kwargs)
776 execution_info = _loop.run_until_complete(_execute(
777 *set(to_execute_tileables), session=self,
--> 778 show_progress=show_progress, **kwargs))
779 if wait:
780 return tileable if len(tileables) == 0 else \

~/miniconda3/lib/python3.7/asyncio/base_events.py in run_until_complete(self, future)
585 raise RuntimeError('Event loop stopped before Future completed.')
586
--> 587 return future.result()
588
589 def stop(self):

~/miniconda3/lib/python3.7/site-packages/mars/core/session.py in _execute(session, wait, show_progress, progress_update_interval, cancelled, *tileables, **kwargs)
635 try:
636 await asyncio.wait_for(asyncio.shield(execution_info),
--> 637 progress_update_interval)
638 # done
639 if not cancelled.is_set():

~/miniconda3/lib/python3.7/asyncio/tasks.py in wait_for(fut, timeout, loop)
440
441 if fut.done():
--> 442 return fut.result()
443 else:
444 fut.remove_done_callback(cb)

~/miniconda3/lib/python3.7/asyncio/tasks.py in _wrap_awaitable(awaitable)
628 that will later be wrapped in a Task by ensure_future().
629 """
--> 630 return (yield from awaitable.__await__())
631
632

~/miniconda3/lib/python3.7/asyncio/events.py in _run(self)
86 def _run(self):
87 try:
---> 88 self._context.run(self._callback, *self._args)
89 except Exception as exc:
90 cb = format_helpers._format_callback_source(

~/miniconda3/lib/python3.7/site-packages/mars/core/session.py in _attach_session(fut)
610
611 def _attach_session(fut: asyncio.Future):
--> 612 fut.result()
613 for t in tileables:
614 t._attach_session(session)

~/miniconda3/lib/python3.7/site-packages/mars/deploy/oscar/session.py in _run_in_background(self, tileables, task_id, progress)
158 await self._task_api.cancel_task(task_id)
159 if task_result.error:
--> 160 raise task_result.error.with_traceback(task_result.traceback)
161 if cancelled:
162 return

~/miniconda3/lib/python3.7/site-packages/mars/services/scheduling/worker/execution.py in internal_run_subtask(self, subtask, band_name)
221 run_aiotask = asyncio.create_task(
222 subtask_api.run_subtask_in_slot(band_name, slot_id, subtask))
--> 223 yield asyncio.shield(run_aiotask)
224 if subtask_info.cancelling:
225 raise asyncio.CancelledError

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor._run_actor_async_generator()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in _handle_actor_result()

~/miniconda3/lib/python3.7/site-packages/mars/services/subtask/api.py in run_subtask_in_slot(self, band_name, slot_id, subtask)
50 """
51 ref = await self._get_runner_ref(band_name, slot_id)
---> 52 return await ref.run_subtask(subtask)
53
54 async def cancel_subtask_in_slot(self, band_name: str, slot_id: int):

~/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py in send(self, actor_ref, message, wait_response)
149 if wait_response:
150 result = await self._wait(future, actor_ref.address, message)
--> 151 return self._process_result_message(result)
152 else:
153 return future

~/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py in _process_result_message(message)
57 return message.result
58 else:
---> 59 raise message.error.with_traceback(message.traceback)
60
61 async def _wait(self,

~/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/pool.py in send()
480 coro = self._actors[actor_id].__on_receive__(message.content)
481 with self._run_coro(message.message_id, coro) as future:
--> 482 result = await future
483 processor.result = ResultMessage(message.message_id, result,
484 protocol=message.protocol)

~/miniconda3/lib/python3.7/site-packages/mars/oscar/api.py in __on_receive__()
109 Message shall be (method_name,) + args + (kwargs,)
110 """
--> 111 return await super().__on_receive__(message)
112
113

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in __on_receive__()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor.__on_receive__()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in _handle_actor_result()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor._run_actor_async_generator()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor._run_actor_async_generator()

~/miniconda3/lib/python3.7/site-packages/mars/services/subtask/worker/runner.py in run_subtask()
83 self._running_processor = self._last_processor = processor
84 try:
---> 85 result = yield self._running_processor.run(subtask)
86 finally:
87 self._running_processor = None

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor._run_actor_async_generator()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in _handle_actor_result()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py in send()
149 if wait_response:
150 result = await self._wait(future, actor_ref.address, message)
--> 151 return self._process_result_message(result)
152 else:
153 return future

~/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py in _process_result_message()
57 return message.result
58 else:
---> 59 raise message.error.with_traceback(message.traceback)
60
61 async def _wait(self,

~/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/pool.py in send()
480 coro = self._actors[actor_id].__on_receive__(message.content)
481 with self._run_coro(message.message_id, coro) as future:
--> 482 result = await future
483 processor.result = ResultMessage(message.message_id, result,
484 protocol=message.protocol)

~/miniconda3/lib/python3.7/site-packages/mars/oscar/api.py in __on_receive__()
109 Message shall be (method_name,) + args + (kwargs,)
110 """
--> 111 return await super().__on_receive__(message)
112
113

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in __on_receive__()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor.__on_receive__()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in _handle_actor_result()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor._run_actor_async_generator()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor._run_actor_async_generator()

~/miniconda3/lib/python3.7/site-packages/mars/services/subtask/worker/processor.py in run()
422 self._running_aio_task = asyncio.create_task(processor.run())
423 try:
--> 424 result = yield self._running_aio_task
425 raise mo.Return(result)
426 finally:

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor._run_actor_async_generator()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in _handle_actor_result()

~/miniconda3/lib/python3.7/site-packages/mars/services/subtask/worker/processor.py in run()
261 put_infos = asyncio.create_task(self._storage_api.put.batch(*puts))
262 try:
--> 263 store_infos = await put_infos
264 store_infos_iter = iter(store_infos)
265 for data_key, puts in data_key_to_puts.items():

~/miniconda3/lib/python3.7/site-packages/mars/utils.py in _async_batch()
1150 if self.batch_func:
1151 args_list, kwargs_list = self._gen_args_kwargs_list(delays)
-> 1152 return await self.batch_func(args_list, kwargs_list)
1153 else:
1154 # this function has no batch implementation

~/miniconda3/lib/python3.7/site-packages/mars/services/storage/api/oscar.py in batch_put()
107 self._session_id, *args, **kwargs)
108 )
--> 109 return await self._storage_handler_ref.put.batch(*puts)
110
111 @extensible

~/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py in send()
149 if wait_response:
150 result = await self._wait(future, actor_ref.address, message)
--> 151 return self._process_result_message(result)
152 else:
153 return future

~/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/context.py in _process_result_message()
57 return message.result
58 else:
---> 59 raise message.error.with_traceback(message.traceback)
60
61 async def _wait(self,

~/miniconda3/lib/python3.7/site-packages/mars/oscar/backends/pool.py in send()
480 coro = self._actors[actor_id].__on_receive__(message.content)
481 with self._run_coro(message.message_id, coro) as future:
--> 482 result = await future
483 processor.result = ResultMessage(message.message_id, result,
484 protocol=message.protocol)

~/miniconda3/lib/python3.7/site-packages/mars/oscar/api.py in __on_receive__()
109 Message shall be (method_name,) + args + (kwargs,)
110 """
--> 111 return await super().__on_receive__(message)
112
113

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in __on_receive__()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor.__on_receive__()

~/miniconda3/lib/python3.7/site-packages/mars/oscar/core.pyx in mars.oscar.core._Actor.__on_receive__()

~/miniconda3/lib/python3.7/site-packages/mars/utils.py in _async_batch()
1150 if self.batch_func:
1151 args_list, kwargs_list = self._gen_args_kwargs_list(delays)
-> 1152 return await self.batch_func(args_list, kwargs_list)
1153 else:
1154 # this function has no batch implementation

~/miniconda3/lib/python3.7/site-packages/mars/services/storage/core.py in batch_put()
334 for size, data_key, obj in zip(sizes, data_keys, objs):
335 logger.debug(f'Begin to put data key {data_key}')
--> 336 object_info = await self._clients[level].put(obj)
337 data_info = _build_data_info(object_info, level, size)
338 data_infos.append(data_info)

~/miniconda3/lib/python3.7/site-packages/mars/storage/shared_memory.py in put()
148
149 serializer = AioSerializer(obj)
--> 150 buffers = await serializer.run()
151 buffer_size = sum(getattr(buf, 'nbytes', len(buf))
152 for buf in buffers)

~/miniconda3/lib/python3.7/site-packages/mars/serialization/aio.py in run()
57
58 async def run(self):
---> 59 return self._get_buffers()
60
61

~/miniconda3/lib/python3.7/site-packages/mars/serialization/aio.py in _get_buffers()
32
33 def _get_buffers(self):
---> 34 headers, buffers = serialize(self._obj)
35
36 # add buffer lengths into headers

~/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py in serialize()
337
338 serializer = _serial_dispatcher.get_handler(type(obj))
--> 339 result = serializer.serialize(obj, context)
340
341 if not isinstance(result, types.GeneratorType):

~/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py in wrapped()
73 else:
74 context[id(obj)] = obj
---> 75 return func(self, obj, context)
76 return wrapped
77

~/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py in serialize()
150 @buffered
151 def serialize(self, obj, context: Dict):
--> 152 return {}, pickle_buffers(obj)
153
154 def deserialize(self, header: Dict, buffers: List, context: Dict):

~/miniconda3/lib/python3.7/site-packages/mars/serialization/core.py in pickle_buffers()
90 obj,
91 buffer_callback=buffer_cb,
---> 92 protocol=BUFFER_PICKLE_PROTOCOL,
93 )
94 else: # pragma: no cover

~/miniconda3/lib/python3.7/site-packages/cloudpickle/cloudpickle_fast.py in dumps()
71 file, protocol=protocol, buffer_callback=buffer_callback
72 )
---> 73 cp.dump(obj)
74 return file.getvalue()
75

~/miniconda3/lib/python3.7/site-packages/cloudpickle/cloudpickle_fast.py in dump()
561 def dump(self, obj):
562 try:
--> 563 return Pickler.dump(self, obj)
564 except RuntimeError as e:
565 if "recursion" in e.args[0]:

TypeError: can't pickle weakref objects
```

Contributor guide

Open the contributing guide

Research direction

Reproduce the shown mars.dataframe.DataFrame.apply call with the mfunc wrapper and inspect the traceback from mars/services/subtask/worker/processor.py through mars/storage/shared_memory.py and mars/serialization/core.py. Trace the object that reaches cloudpickle and verify that the same operation completes without the reported weakref pickling failure.

Written by the indexing model from the issue text.

Assessment

Tech stack
python
Domain
data
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.