From 2e0346ff9cd741a9e493d50e1ea531656df1a9f6 Mon Sep 17 00:00:00 2001 From: Gaetan Lepage Date: Mon, 3 Aug 2026 11:33:29 +0000 Subject: [PATCH] python3Packages.apache-beam: relax numpy, fix pandas 3 compat --- .../python-modules/apache-beam/default.nix | 39 ++++ .../apache-beam/pandas-3-compat.patch | 169 ++++++++++++++++++ 2 files changed, 208 insertions(+) create mode 100644 pkgs/development/python-modules/apache-beam/pandas-3-compat.patch diff --git a/pkgs/development/python-modules/apache-beam/default.nix b/pkgs/development/python-modules/apache-beam/default.nix index a3c834e0e751..7675694f7d25 100644 --- a/pkgs/development/python-modules/apache-beam/default.nix +++ b/pkgs/development/python-modules/apache-beam/default.nix @@ -79,13 +79,23 @@ buildPythonPackage (finalAttrs: { sourceRoot = "${finalAttrs.src.name}/sdks/python"; + patches = [ + # Beam's dataframe module references DataFrame.first/last/bool/applymap, all removed in + # pandas 3. Their `with_docs_from` decorators do getattr at import, crashing frames.py (which + # then never registers the deferred types). Guard them behind hasattr like the neighboring + # ffill/bfill aliases. Upstream caps pandas < 2.4, so this is nixpkgs-local. + ./pandas-3-compat.patch + ]; + postPatch = '' substituteInPlace pyproject.toml \ --replace-fail "distlib==0.4.2" "distlib" \ --replace-fail "cython>=3.2.5,<4" "cython" \ + --replace-fail "numpy>=1.14.3,<2.5.0" "numpy" \ --replace-fail "==" ">=" substituteInPlace setup.py \ + --replace-fail "numpy>=1.14.3,<2.5.0" "numpy" \ --replace-fail " copy_tests_from_docs()" "" ''; @@ -94,6 +104,7 @@ buildPythonPackage (finalAttrs: { "envoy-data-plane" "httplib2" "jsonpickle" + "numpy" "objsize" "protobuf" "pyarrow" @@ -238,6 +249,12 @@ buildPythonPackage (finalAttrs: { # grpc_status:14, grpc_message:"Cancelling all calls"}" # Upstream issue https://github.com/apache/beam/issues/33851 "apache_beam/runners/portability/portable_runner_test.py" + + # Pandas 3 behavioral changes in Beam's dataframe module (upstream caps pandas < 2.4): + # the doctests compare against pandas-version-specific output and the taxiride example + # relies on DataFrame APIs removed in pandas 3. + "apache_beam/dataframe/doctests_test.py" + "apache_beam/examples/dataframe/taxiride_test.py" ] ++ lib.optionals (pythonAtLeast "3.13") [ # > instruction = ofs_table[pc] @@ -307,6 +324,28 @@ buildPythonPackage (finalAttrs: { # AssertionError: False is not true "test_samples_all_with_both_experiments" + + # ValueError: buffer source array is read-only (Cython coder over a read-only numpy buffer) + "test_coders_microbenchmark" + + # Pandas 3 behavioral changes in Beam's dataframe module (upstream caps pandas < 2.4). + # groupby() no longer accepts the removed `axis` argument: + "test_groupby_apply" + "test_groupby_sum_mean" + "test_scalar" + # (only surface on Python 3.13, but same groupby() `axis` root cause) + "test_dataframes_with_grouped_index" + "test_dataframes_with_multi_index" + "test_dataframes_with_multi_index_get_result" + # read_csv attribute/dtype differences: + "test_csv_splitter_00_defaults" + "test_csv_splitter_01_header" + "test_csv_splitter_05_names_and_header" + "test_csv_splitter_06_skip_blank_lines" + "test_csv_splitter_07_skip_blank_lines" + "test_csv_splitter_08_comment" + # nan instead of None for missing values: + "test_convert_with_none" ] ++ lib.optionals stdenv.hostPlatform.isDarwin [ # PermissionError: [Errno 13] Permission denied: '/tmp/...' diff --git a/pkgs/development/python-modules/apache-beam/pandas-3-compat.patch b/pkgs/development/python-modules/apache-beam/pandas-3-compat.patch new file mode 100644 index 000000000000..b15907528659 --- /dev/null +++ b/pkgs/development/python-modules/apache-beam/pandas-3-compat.patch @@ -0,0 +1,169 @@ +--- 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)