This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/paimon.git
The following commit(s) were added to refs/heads/master by this push:
new 3ed6fa6769 [python] Introduce theta_sketch aggregator function (#8760)
3ed6fa6769 is described below
commit 3ed6fa676986567c81a38a7b0bcba2f99a3c67c1
Author: AuroraVoyage <[email protected]>
AuthorDate: Sun Jul 26 19:44:46 2026 +0800
[python] Introduce theta_sketch aggregator function (#8760)
---
.github/workflows/paimon-python-checks.yml | 8 ++--
paimon-python/dev/requirements.txt | 2 +
.../pypaimon/read/merge_engine_support.py | 1 +
.../pypaimon/read/reader/aggregate/aggregators.py | 45 ++++++++++++++++++++++
.../pypaimon/tests/test_field_aggregators.py | 29 ++++++++++++++
5 files changed, 81 insertions(+), 4 deletions(-)
diff --git a/.github/workflows/paimon-python-checks.yml
b/.github/workflows/paimon-python-checks.yml
index efa187367b..00540533f1 100755
--- a/.github/workflows/paimon-python-checks.yml
+++ b/.github/workflows/paimon-python-checks.yml
@@ -119,7 +119,7 @@ jobs:
if [[ "${{ matrix.python-version }}" == "3.6.15" ]]; then
python -m pip install --upgrade pip==21.3.1
python --version
- python -m pip install --no-cache-dir pyroaring
readerwriterlock==1.0.9 'fsspec==2021.10.1' 'cachetools==4.2.4'
'ossfs==2021.8.0' pyarrow==6.0.1 pandas==1.1.5 'polars==0.9.12'
'fastavro==1.4.7' zstandard==0.19.0 dataclasses==0.8.0 flake8 pytest
py4j==0.10.9.9 requests parameterized==0.8.1 2>&1 >/dev/null
+ python -m pip install --no-cache-dir pyroaring
readerwriterlock==1.0.9 'fsspec==2021.10.1' 'cachetools==4.2.4'
'ossfs==2021.8.0' pyarrow==6.0.1 pandas==1.1.5 'polars==0.9.12'
'fastavro==1.4.7' zstandard==0.19.0 dataclasses==0.8.0 flake8 pytest
py4j==0.10.9.9 requests parameterized==0.8.1 datasketches==4.1.0 2>&1 >/dev/null
python -m pip install 'lumina-data>=${{ env.LUMINA_DATA_VERSION
}}' -i https://pypi.org/simple/
elif [[ "${{ matrix.python-version }}" == "3.7" ]]; then
# 3.7 installs the version-pinned set declared for 3.7 in
dev/requirements.txt.
@@ -131,7 +131,7 @@ jobs:
else
python -m pip install --upgrade pip
pip install torch --index-url https://download.pytorch.org/whl/cpu
- python -m pip install pyroaring readerwriterlock==1.0.9
fsspec==2024.3.1 cachetools==5.3.3 ossfs==2023.12.0 ray==2.54.0
fastavro==1.11.1 'isal>=1.8,<2' pyarrow==16.0.0 zstandard==0.24.0
polars==1.32.0 duckdb==1.3.2 numpy==1.24.3 pandas==2.0.3 pylance==0.39.0
cramjam flake8==4.0.1 pytest~=7.0 py4j==0.10.9.9 requests parameterized==0.9.0
'daft>=0.7.6' 'datafusion>=52'
+ python -m pip install pyroaring readerwriterlock==1.0.9
fsspec==2024.3.1 cachetools==5.3.3 ossfs==2023.12.0 ray==2.54.0
fastavro==1.11.1 'isal>=1.8,<2' pyarrow==16.0.0 zstandard==0.24.0
polars==1.32.0 duckdb==1.3.2 numpy==1.24.3 pandas==2.0.3 pylance==0.39.0
cramjam flake8==4.0.1 pytest~=7.0 py4j==0.10.9.9 requests parameterized==0.9.0
'daft>=0.7.6' 'datafusion>=52' datasketches
if [[ "${{ matrix.python-version }}" == "3.11" ]]; then
# Exercise the split-planning API in one lane until the
compatible 0.3 wheel is published.
python -m pip install
"git+https://github.com/apache/paimon-rust.git@${PYPAIMON_RUST_REV}#subdirectory=bindings/python"
@@ -179,7 +179,7 @@ jobs:
run: |
python -m pip install --upgrade pip
pip install torch --index-url https://download.pytorch.org/whl/cpu
- python -m pip install pyroaring readerwriterlock==1.0.9
fsspec==2024.3.1 cachetools==5.3.3 ossfs==2023.12.0 ray==2.54.0
fastavro==1.11.1 pyarrow==16.0.0 zstandard==0.24.0 polars==1.32.0 duckdb==1.3.2
numpy==1.24.3 pandas==2.0.3 pylance==0.39.0 flake8==4.0.1 pytest~=7.0
py4j==0.10.9.9 requests parameterized==0.9.0
+ python -m pip install pyroaring readerwriterlock==1.0.9
fsspec==2024.3.1 cachetools==5.3.3 ossfs==2023.12.0 ray==2.54.0
fastavro==1.11.1 pyarrow==16.0.0 zstandard==0.24.0 polars==1.32.0 duckdb==1.3.2
numpy==1.24.3 pandas==2.0.3 pylance==0.39.0 flake8==4.0.1 pytest~=7.0
py4j==0.10.9.9 requests parameterized==0.9.0 datasketches
python -m pip install 'lumina-data>=${{ env.LUMINA_DATA_VERSION
}}' -i https://pypi.org/simple/
- name: Run lint-python.sh
shell: bash
@@ -278,7 +278,7 @@ jobs:
pyroaring readerwriterlock==1.0.9 fsspec==2024.3.1
cachetools==5.3.3 ossfs==2023.12.0 \
fastavro==1.11.1 pyarrow==16.0.0 zstandard==0.24.0 polars==1.32.0
duckdb==1.3.2 \
numpy==1.24.3 pandas==2.0.3 cramjam pytest~=7.0 py4j==0.10.9.9
requests \
- parameterized==0.9.0 packaging
+ parameterized==0.9.0 packaging datasketches
python -m pip install 'lumina-data>=${{ env.LUMINA_DATA_VERSION }}'
-i https://pypi.org/simple/
- name: Test Ray version compatibility
run: |
diff --git a/paimon-python/dev/requirements.txt
b/paimon-python/dev/requirements.txt
index dc9b4b4e91..e2697f0b39 100644
--- a/paimon-python/dev/requirements.txt
+++ b/paimon-python/dev/requirements.txt
@@ -42,3 +42,5 @@ zstandard>=0.19,<1
backports.zstd>=1.0.0,<1.4.0; python_version >= "3.9" and python_version <
"3.14"
cramjam>=1.3.0,<3; python_version>="3.7"
pyyaml>=5.4,<7
+datasketches>=4,<5; python_version < "3.9"
+datasketches>=5,<6; python_version >= "3.9"
diff --git a/paimon-python/pypaimon/read/merge_engine_support.py
b/paimon-python/pypaimon/read/merge_engine_support.py
index 321421b892..d6b38c20c6 100644
--- a/paimon-python/pypaimon/read/merge_engine_support.py
+++ b/paimon-python/pypaimon/read/merge_engine_support.py
@@ -67,6 +67,7 @@ _AGGREGATION_SUPPORTED_AGG_FUNCS = frozenset([
"collect",
"merge_map_with_keytime",
"merge_map",
+ "theta_sketch",
])
_FIELDS_PREFIX = "fields."
_FIELD_SEQUENCE_GROUP_SUFFIX = ".sequence-group"
diff --git a/paimon-python/pypaimon/read/reader/aggregate/aggregators.py
b/paimon-python/pypaimon/read/reader/aggregate/aggregators.py
index 5c325d99eb..4d06e303f9 100644
--- a/paimon-python/pypaimon/read/reader/aggregate/aggregators.py
+++ b/paimon-python/pypaimon/read/reader/aggregate/aggregators.py
@@ -35,6 +35,8 @@ error rather than a silent fallback.
from typing import Any, List, Dict, Optional, Tuple, Union, Set
+from _datasketches import compact_theta_sketch, theta_union
+
from pypaimon.common.options import CoreOptions
from pypaimon.common.options.core_options import NestedKeyNullStrategy
from pypaimon.read.reader.aggregate import register_aggregator
@@ -65,6 +67,7 @@ NAME_NESTED_PARTIAL_UPDATE = "nested_partial_update"
NAME_COLLECT = "collect"
NAME_MERGE_MAP_WITH_KEYTIME = "merge_map_with_keytime"
NAME_MERGE_MAP = "merge_map"
+NAME_THETA_SKETCH = "theta_sketch"
# Base SQL type names treated as numeric for sum/product-style
@@ -1063,6 +1066,45 @@ class FieldMergeMapAgg(FieldAggregator):
return keys
+class FieldThetaSketchAgg(FieldAggregator):
+ """Aggregator for ThetaSketch."""
+
+ def __init__(self, name: str, field_type: DataType):
+ super().__init__(name, field_type)
+ if _atomic_base_name(field_type) not in ("VARBINARY", "BYTES"):
+ raise ValueError(
+ "Data type for theta sketch column must be 'VarBinaryType' but
was '{}'.".format(field_type)
+ )
+
+ def agg(self, accumulator: Any, input_field: Any) -> Any:
+ if accumulator is None or input_field is None:
+ return input_field if accumulator is None else accumulator
+
+ if not isinstance(accumulator, (bytes, bytearray)):
+ raise TypeError(
+ "ThetaSketch accumulator must be bytes, got
{}".format(type(accumulator))
+ )
+
+ if not isinstance(input_field, (bytes, bytearray)):
+ raise TypeError(
+ "ThetaSketch input must be bytes, got
{}".format(type(input_field))
+ )
+
+ if isinstance(accumulator, bytearray):
+ accumulator = bytes(accumulator)
+ if isinstance(input_field, bytearray):
+ input_field = bytes(input_field)
+
+ sketch1 = compact_theta_sketch.deserialize(accumulator)
+ sketch2 = compact_theta_sketch.deserialize(input_field)
+
+ union = theta_union()
+ union.update(sketch1)
+ union.update(sketch2)
+
+ return union.get_result().serialize()
+
+
# ---------------------------------------------------------------------------
# Registration. Each builder binds an identifier to a factory that
# optionally validates the column DataType before constructing the
@@ -1151,3 +1193,6 @@ register_aggregator(
register_aggregator(
NAME_MERGE_MAP, _build_no_type_check(FieldMergeMapAgg, NAME_MERGE_MAP)
)
+register_aggregator(
+ NAME_THETA_SKETCH, _build_no_type_check(FieldThetaSketchAgg,
NAME_THETA_SKETCH)
+)
diff --git a/paimon-python/pypaimon/tests/test_field_aggregators.py
b/paimon-python/pypaimon/tests/test_field_aggregators.py
index 107b160df1..f04da32a11 100644
--- a/paimon-python/pypaimon/tests/test_field_aggregators.py
+++ b/paimon-python/pypaimon/tests/test_field_aggregators.py
@@ -30,6 +30,8 @@ from decimal import Decimal
from functools import reduce
from typing import List
+from _datasketches import update_theta_sketch
+
from pypaimon.common.options import CoreOptions, Options
from pypaimon.data import Timestamp
from pypaimon.read.reader.aggregate import create_field_aggregator
@@ -50,6 +52,7 @@ from pypaimon.read.reader.aggregate.aggregators import (
FieldCollectAgg,
FieldMergeMapWithKeyTimeAgg,
FieldMergeMapAgg,
+ FieldThetaSketchAgg,
)
from pypaimon.schema.data_types import AtomicType, DataField, RowType,
ArrayType, MapType
from pypaimon.table.row.generic_row import GenericRow
@@ -2106,6 +2109,32 @@ class FieldMergeMapAggTest(unittest.TestCase):
self.assertEqual(acc, {2: "B"})
+class FieldThetaSketchAggTest(unittest.TestCase):
+ @staticmethod
+ def sketch_of(*values: int) -> bytes:
+ sketch = update_theta_sketch()
+
+ for value in values:
+ sketch.update(value)
+
+ return sketch.compact().serialize()
+
+ def test_field_theta_sketch_agg(self):
+ agg = _make("theta_sketch", "VARBINARY(20)")
+ self.assertIsInstance(agg, FieldThetaSketchAgg)
+
+ input_val = self.sketch_of(1)
+ acc1 = self.sketch_of(2, 3)
+ acc2 = self.sketch_of(1, 2, 3)
+
+ self.assertIsNone(agg.agg(None, None))
+
+ self.assertEqual(agg.agg(None, input_val), input_val)
+ self.assertEqual(agg.agg(acc1, None), acc1)
+ self.assertEqual(agg.agg(acc1, input_val), acc2)
+ self.assertEqual(agg.agg(acc2, input_val), acc2)
+
+
class RegistrationTest(unittest.TestCase):
"""Sanity check that all 10 expected aggregators (the primary-key
placeholder plus 9 value aggregators) are registered when the