MCPcopy Create free account
hub / github.com/dask/dask / P2PShuffle

Class P2PShuffle

dask/dataframe/dask_expr/_shuffle.py:555–625  ·  view source on GitHub ↗

P2P worker-based shuffle implementation

Source from the content-addressed store, hash-verified

553
554
555class P2PShuffle(SimpleShuffle):
556 """P2P worker-based shuffle implementation"""
557
558 @functools.cached_property
559 def _meta(self):
560 return self.frame._meta.drop(columns=self.partitioning_index)
561
562 def _layer(self):
563 from distributed.shuffle._core import (
564 P2PBarrierTask,
565 ShuffleId,
566 barrier_key,
567 p2p_barrier,
568 )
569 from distributed.shuffle._shuffle import DataFrameShuffleSpec, shuffle_unpack
570
571 dsk = {}
572 token = self._name.split("-")[-1]
573 shuffle_id = ShuffleId(token)
574 _barrier_key = barrier_key(shuffle_id)
575 name = f"shuffle-transfer-{token}"
576
577 parts_out = (
578 self._partitions if self._filtered else list(range(self.npartitions_out))
579 )
580 # Avoid embedding a materialized list unless necessary
581 parts_out_arg = (
582 tuple(self._partitions) if self._filtered else self.npartitions_out
583 )
584
585 transfer_keys = list()
586 for i in range(self.frame.npartitions):
587 t = Task(
588 (name, i),
589 _shuffle_transfer,
590 TaskRef((self.frame._name, i)),
591 token,
592 i,
593 )
594 dsk[t.key] = t
595 transfer_keys.append(t.ref())
596
597 barrier = P2PBarrierTask(
598 _barrier_key,
599 p2p_barrier,
600 token,
601 *transfer_keys,
602 spec=DataFrameShuffleSpec(
603 id=shuffle_id,
604 npartitions=self.npartitions_out,
605 column=self.partitioning_index,
606 meta=self.frame._meta,
607 parts_out=parts_out_arg,
608 disk=True,
609 drop_column=True,
610 ),
611 )
612 dsk[barrier.key] = barrier

Callers 1

_lowerMethod · 0.85

Calls

no outgoing calls

Tested by

no test coverage detected