--- a/apache_beam/dataframe/frames.py +++ b/apache_beam/dataframe/frames.py @@ -335,35 +335,37 @@ if hasattr(pd.DataFrame, 'pad'): pad = _fillna_alias('pad') - @frame_base.with_docs_from(pd.DataFrame) - def first(self, offset): - per_partition = expressions.ComputedExpression( - 'first-per-partition', lambda df: df.sort_index().first(offset=offset), - [self._expr], - preserves_partition_by=partitionings.Arbitrary(), - requires_partition_by=partitionings.Arbitrary()) - with expressions.allow_non_parallel_operations(True): - return frame_base.DeferredFrame.wrap( - expressions.ComputedExpression( - 'first', lambda df: df.sort_index().first(offset=offset), - [per_partition], - preserves_partition_by=partitionings.Arbitrary(), - requires_partition_by=partitionings.Singleton())) + if hasattr(pd.DataFrame, 'first'): + @frame_base.with_docs_from(pd.DataFrame) + def first(self, offset): + per_partition = expressions.ComputedExpression( + 'first-per-partition', lambda df: df.sort_index().first(offset=offset), + [self._expr], + preserves_partition_by=partitionings.Arbitrary(), + requires_partition_by=partitionings.Arbitrary()) + with expressions.allow_non_parallel_operations(True): + return frame_base.DeferredFrame.wrap( + expressions.ComputedExpression( + 'first', lambda df: df.sort_index().first(offset=offset), + [per_partition], + preserves_partition_by=partitionings.Arbitrary(), + requires_partition_by=partitionings.Singleton())) - @frame_base.with_docs_from(pd.DataFrame) - def last(self, offset): - per_partition = expressions.ComputedExpression( - 'last-per-partition', lambda df: df.sort_index().last(offset=offset), - [self._expr], - preserves_partition_by=partitionings.Arbitrary(), - requires_partition_by=partitionings.Arbitrary()) - with expressions.allow_non_parallel_operations(True): - return frame_base.DeferredFrame.wrap( - expressions.ComputedExpression( - 'last', lambda df: df.sort_index().last(offset=offset), - [per_partition], - preserves_partition_by=partitionings.Arbitrary(), - requires_partition_by=partitionings.Singleton())) + if hasattr(pd.DataFrame, 'last'): + @frame_base.with_docs_from(pd.DataFrame) + def last(self, offset): + per_partition = expressions.ComputedExpression( + 'last-per-partition', lambda df: df.sort_index().last(offset=offset), + [self._expr], + preserves_partition_by=partitionings.Arbitrary(), + requires_partition_by=partitionings.Arbitrary()) + with expressions.allow_non_parallel_operations(True): + return frame_base.DeferredFrame.wrap( + expressions.ComputedExpression( + 'last', lambda df: df.sort_index().last(offset=offset), + [per_partition], + preserves_partition_by=partitionings.Arbitrary(), + requires_partition_by=partitionings.Singleton())) @frame_base.with_docs_from(pd.DataFrame) @frame_base.args_to_kwargs(pd.DataFrame) @@ -798,27 +800,28 @@ requires_partition_by=partitionings.Singleton(), preserves_partition_by=partitionings.Singleton())) - @frame_base.with_docs_from(pd.DataFrame) - def bool(self): - # TODO: Documentation about DeferredScalar - # Will throw if any partition has >1 element - bools = expressions.ComputedExpression( - 'get_bools', - # Wrap scalar results in a Series for easier concatenation later - lambda df: pd.Series([], dtype=bool) - if df.empty else pd.Series([df.bool()]), - [self._expr], - requires_partition_by=partitionings.Arbitrary(), - preserves_partition_by=partitionings.Singleton()) + if hasattr(pd.DataFrame, 'bool'): + @frame_base.with_docs_from(pd.DataFrame) + def bool(self): + # TODO: Documentation about DeferredScalar + # Will throw if any partition has >1 element + bools = expressions.ComputedExpression( + 'get_bools', + # Wrap scalar results in a Series for easier concatenation later + lambda df: pd.Series([], dtype=bool) + if df.empty else pd.Series([df.bool()]), + [self._expr], + requires_partition_by=partitionings.Arbitrary(), + preserves_partition_by=partitionings.Singleton()) - with expressions.allow_non_parallel_operations(True): - # Will throw if overall dataset has != 1 element - return frame_base.DeferredFrame.wrap( - expressions.ComputedExpression( - 'combine_all_bools', lambda bools: bools.bool(), [bools], - proxy=bool(), - requires_partition_by=partitionings.Singleton(), - preserves_partition_by=partitionings.Singleton())) + with expressions.allow_non_parallel_operations(True): + # Will throw if overall dataset has != 1 element + return frame_base.DeferredFrame.wrap( + expressions.ComputedExpression( + 'combine_all_bools', lambda bools: bools.bool(), [bools], + proxy=bool(), + requires_partition_by=partitionings.Singleton(), + preserves_partition_by=partitionings.Singleton())) @frame_base.with_docs_from(pd.DataFrame) def equals(self, other): @@ -2932,7 +2935,8 @@ agg = aggregate - applymap = frame_base._elementwise_method('applymap', base=pd.DataFrame) + if hasattr(pd.DataFrame, 'applymap'): + applymap = frame_base._elementwise_method('applymap', base=pd.DataFrame) if PD_VERSION >= (2, 1): map = frame_base._elementwise_method('map', base=pd.DataFrame) add_prefix = frame_base._elementwise_method('add_prefix', base=pd.DataFrame) @@ -4414,18 +4418,19 @@ return self.apply(apply_fn).droplevel(self._grouping_columns) - @property # type: ignore - @frame_base.with_docs_from(DataFrameGroupBy) - def dtypes(self): - return frame_base.DeferredFrame.wrap( - expressions.ComputedExpression( - 'dtypes', - lambda gb: gb.dtypes, - [self._expr], - requires_partition_by=partitionings.Arbitrary(), - preserves_partition_by=partitionings.Arbitrary() - ) - ) + if hasattr(DataFrameGroupBy, 'dtypes'): + @property # type: ignore + @frame_base.with_docs_from(DataFrameGroupBy) + def dtypes(self): + return frame_base.DeferredFrame.wrap( + expressions.ComputedExpression( + 'dtypes', + lambda gb: gb.dtypes, + [self._expr], + requires_partition_by=partitionings.Arbitrary(), + preserves_partition_by=partitionings.Arbitrary() + ) + ) if hasattr(DataFrameGroupBy, 'value_counts'): @frame_base.with_docs_from(DataFrameGroupBy) @@ -4745,7 +4750,8 @@ describe = frame_base.not_implemented_method('describe', base_type=DataFrameGroupBy) diff = frame_base._elementwise_method('diff', base=DataFrameGroupBy) - fillna = frame_base._elementwise_method('fillna', base=DataFrameGroupBy) + if hasattr(DataFrameGroupBy, 'fillna'): + fillna = frame_base._elementwise_method('fillna', base=DataFrameGroupBy) filter = frame_base._elementwise_method('filter', base=DataFrameGroupBy) first = frame_base._elementwise_method('first', base=DataFrameGroupBy) get_group = frame_base._elementwise_method('get_group', base=DataFrameGroupBy)