Skip to content

Commit

Permalink
dask repartition (#12)
Browse files Browse the repository at this point in the history
  • Loading branch information
Han Wang committed May 11, 2020
1 parent 86e2f58 commit 87fd3c4
Show file tree
Hide file tree
Showing 3 changed files with 4 additions and 3 deletions.
1 change: 1 addition & 0 deletions fugue_dask/execution_engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,7 @@ def _map(pdf: Any) -> pd.DataFrame:
pdf = self.repartition(df, partition_spec)
result = pdf.native.map_partitions(_map, meta=output_schema.pandas_dtype)
else:
df = self.repartition(df, PartitionSpec(num=partition_spec.num_partitions))
result = DASK_UTILS.safe_groupby_apply(
df.native,
partition_spec.partition_by,
Expand Down
4 changes: 2 additions & 2 deletions fugue_test/builtin_suite.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,8 +33,8 @@ def dag(self) -> "DagTester":

def test_create_show(self):
with self.dag() as dag:
dag.df([[0]], "a:int").persist(123).partition(num=2).show()
dag.df(ArrayDataFrame([[0]], "a:int")).persist(456).broadcast().show()
dag.df([[0]], "a:int").persist().partition(num=2).show()
dag.df(ArrayDataFrame([[0]], "a:int")).persist().broadcast().show()

def test_transform(self):
with self.dag() as dag:
Expand Down
2 changes: 1 addition & 1 deletion setup.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
from setuptools import setup, find_packages

VERSION = "0.1.5"
VERSION = "0.1.6"

with open("README.md") as f:
LONG_DESCRIPTION = f.read()
Expand Down

0 comments on commit 87fd3c4

Please sign in to comment.