diff --git a/pemi/pipes/pd.py b/pemi/pipes/pd.py index b75fd0c..60456a1 100644 --- a/pemi/pipes/pd.py +++ b/pemi/pipes/pd.py @@ -175,13 +175,14 @@ def __init__(self, field, forks, **kwargs): ) def flow(self): - grouped = self.sources['main'].df.groupby(self.field) + work_df = self.sources['main'].df.copy() + grouped = work_df.groupby(self.field) for fork in self.forks: if fork in grouped.groups: self.targets[fork].df = grouped.get_group(fork).copy() else: - self.targets[fork].df = pd.DataFrame(columns=self.sources['main'].df.columns) + self.targets[fork].df = pd.DataFrame(columns=work_df.columns) remainder = set(grouped.groups.keys()) - set(self.forks) if len(remainder) > 0: @@ -189,7 +190,13 @@ def flow(self): [grouped.get_group(r) for r in remainder] ).sort_index() else: - self.targets['remainder'].df = pd.DataFrame(columns=self.sources['main'].df.columns) + self.targets['remainder'].df = pd.DataFrame(columns=work_df.columns) + + if None in work_df[self.field].unique(): + self.targets['remainder'].df = pd.concat( + [self.targets['remainder'].df, work_df[work_df[self.field].isna()]] + ).sort_index() + class PdLambdaPipe(pemi.Pipe): ''' diff --git a/tests/pipes/test_pd.py b/tests/pipes/test_pd.py index c59229f..bddd75b 100644 --- a/tests/pipes/test_pd.py +++ b/tests/pipes/test_pd.py @@ -549,8 +549,8 @@ def pipe(self): ) df = pd.DataFrame({ - 'target': ['create', 'update', 'update', 'else1', 'create', 'else2'], - 'values': [1, 2, 3, 4, 5, 6] + 'target': ['create', 'update', 'update', 'else1', 'create', 'else2', '', None], + 'values': [1, 2, 3, 4, 5, 6, 7, 8] }) pipe.sources['main'].df = df @@ -592,9 +592,9 @@ def test_it_forks_data_to_update(self, pipe): def test_it_puts_unknown_values_in_remainder(self, pipe): expected_df = pd.DataFrame({ - 'target': ['else1', 'else2'], - 'values': [4, 6] - }, index=[3, 5]) + 'target': ['else1', 'else2', '', None], + 'values': [4, 6, 7, 8] + }, index=[3, 5, 6, 7]) actual_df = pipe.targets['remainder'].df pt.assert_frame_equal(actual_df, expected_df)