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

Method _layer

dask/dataframe/dask_expr/_shuffle.py:562–625  ·  view source on GitHub ↗
(self)

Source from the content-addressed store, hash-verified

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 = "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
613
614 # TODO: Decompose p2p Into transfer/barrier + unpack
615 name = self._name
616 for i, part_out in enumerate(parts_out):
617 t = Task(
618 (name, i),
619 shuffle_unpack,

Callers 1

_layerMethod · 0.45

Calls 4

TaskClass · 0.90
TaskRefClass · 0.90
splitMethod · 0.80
refMethod · 0.80

Tested by

no test coverage detected