| 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, |