mars-project / mars-project/mars
[BUG]failed call mars.tensor.sum in mars.dataframe.apply
- 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
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