From f673605e2bafde83f07f8e04b75a5acf738cdaa7 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Sat, 10 Oct 2026 12:06:55 +0800 Subject: [PATCH 1/2] [python] Use native builders for local sorted global indexes --- .../globalindex/create_global_index.py | 7 +- .../globalindex/native_index_build.py | 49 ++ .../index/data_evolution_index_source_meta.py | 48 ++ .../tests/native_sorted_index_build_test.py | 488 ++++++++++++++++++ .../write/commit/conflict_detection.py | 6 + .../write/commit/global_index_source_check.py | 82 +++ 6 files changed, 679 insertions(+), 1 deletion(-) create mode 100644 paimon-python/pypaimon/globalindex/native_index_build.py create mode 100644 paimon-python/pypaimon/index/data_evolution_index_source_meta.py create mode 100644 paimon-python/pypaimon/tests/native_sorted_index_build_test.py create mode 100644 paimon-python/pypaimon/write/commit/global_index_source_check.py diff --git a/paimon-python/pypaimon/globalindex/create_global_index.py b/paimon-python/pypaimon/globalindex/create_global_index.py index 569f2f8861f4..ec7fc24488ba 100644 --- a/paimon-python/pypaimon/globalindex/create_global_index.py +++ b/paimon-python/pypaimon/globalindex/create_global_index.py @@ -173,8 +173,13 @@ def build(self, *, execution="local", concurrency=None, ray_remote_args=None) -> from pypaimon.ray.vector_index_build import validate_build_options concurrency, ray_remote_args = validate_build_options( self, concurrency, ray_remote_args) - read_builder = self._table.new_read_builder() partition_filter = self._resolve_partition_filter() + if execution == "local" and self._index_type in _SORTED_INDEX_IDENTIFIERS: + from pypaimon.globalindex.native_index_build import build_native_sorted_index + messages = build_native_sorted_index(self, partition_filter) + if messages is not None: + return messages + read_builder = self._table.new_read_builder() if partition_filter is not None: read_builder = read_builder.with_partition_filter(partition_filter) diff --git a/paimon-python/pypaimon/globalindex/native_index_build.py b/paimon-python/pypaimon/globalindex/native_index_build.py new file mode 100644 index 000000000000..266ebac6eec4 --- /dev/null +++ b/paimon-python/pypaimon/globalindex/native_index_build.py @@ -0,0 +1,49 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""Local sorted-index builds delegated to Rust, retaining Python's commit API.""" + +from pypaimon.read.native_plan import ( + _native_table, _option_value_to_string, _predicate_to_native, native_method_available) +from pypaimon.write.native_commit import _rest_catalog_supported, from_native_commit_messages + + +def build_native_sorted_index(builder, partition_filter): + """Return unpublished messages, or None before building for an ineligible table. + + Rust owns snapshot selection, range planning, sorting and file generation. + Build errors propagate: a failed native build must not start a second build. + """ + table = builder._table + if (builder._index_type not in ('btree', 'bitmap') + or not table.options.native_write_enabled() + or not table.options.data_evolution_enabled() + or not table.options.global_index_enabled() + or table.options.deletion_vectors_enabled() + or table.options.query_auth_enabled + or table.is_primary_key_table + or table.options.file_format() != 'parquet' + or not native_method_available('Table', 'new_sorted_global_index_build_builder') + or not _rest_catalog_supported(table)): + return None + native = _native_table(table).new_sorted_global_index_build_builder() + native.with_index_column(builder._index_columns[0]).with_index_type(builder._index_type) + native.with_options({str(key): _option_value_to_string(value) + for key, value in builder._user_options.items() if value is not None}) + if partition_filter is not None: + native.with_partition_filter(_predicate_to_native(partition_filter)) + return from_native_commit_messages(table, native.build()) diff --git a/paimon-python/pypaimon/index/data_evolution_index_source_meta.py b/paimon-python/pypaimon/index/data_evolution_index_source_meta.py new file mode 100644 index 000000000000..5bdd7d77f08d --- /dev/null +++ b/paimon-python/pypaimon/index/data_evolution_index_source_meta.py @@ -0,0 +1,48 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Java DataEvolutionIndexSourceMeta stored in GlobalIndexMeta.source_meta.""" + +import struct +from dataclasses import dataclass + + +@dataclass(frozen=True) +class DataEvolutionIndexSourceMeta: + scan_snapshot_id: int + + MAGIC = 0x44454958 + VERSION = 1 + + def __post_init__(self): + if self.scan_snapshot_id <= 0: + raise ValueError('Scan snapshot id must be positive.') + + def serialize(self): + return struct.pack('>iiq', self.MAGIC, self.VERSION, self.scan_snapshot_id) + + @classmethod + def is_data_evolution_meta(cls, data): + return data is not None and len(data) >= 4 and data[:4] == b'DEIX' + + @classmethod + def deserialize(cls, data): + if len(data) != 16 or not cls.is_data_evolution_meta(data): + raise ValueError('Invalid data-evolution index source metadata.') + _, version, snapshot_id = struct.unpack('>iiq', data) + if version != cls.VERSION: + raise ValueError('Unsupported data-evolution index source version: {}.'.format(version)) + return cls(snapshot_id) diff --git a/paimon-python/pypaimon/tests/native_sorted_index_build_test.py b/paimon-python/pypaimon/tests/native_sorted_index_build_test.py new file mode 100644 index 000000000000..b56ec6099297 --- /dev/null +++ b/paimon-python/pypaimon/tests/native_sorted_index_build_test.py @@ -0,0 +1,488 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +"""Build Rust index files on REST tables and consume them through Python APIs.""" + +from datetime import date, datetime, time, timezone +from decimal import Decimal +import struct +from unittest.mock import patch + +import pyarrow as pa +import pytest + +from pypaimon import Schema +from pypaimon.api.api_response import ErrorResponse, GetTableSnapshotResponse +from pypaimon.globalindex.create_global_index import GlobalIndexBuilder +from pypaimon.globalindex.data_evolution_global_index_scanner import DataEvolutionGlobalIndexScanner +from pypaimon.index.index_file_handler import IndexFileHandler +from pypaimon.read.native_plan import _native_table, _predicate_to_native, native_method_available +from pypaimon.snapshot.table_snapshot import TableSnapshot +from pypaimon.tests import native_plan_rest_test +from pypaimon.write.native_commit import from_native_commit_messages +from pypaimon.write.table_commit import BatchTableCommit + +rest_catalog = native_plan_rest_test.rest_catalog +pytestmark = [pytest.mark.native_plan, pytest.mark.skipif( + not native_method_available('Table', 'new_sorted_global_index_build_builder'), + reason='Rust index builder required')] + +SCHEMA = pa.schema([('id', pa.int32()), ('name', pa.string()), ('pt', pa.int32())]) +ROWS = [ + {'id': 3, 'name': 'c', 'pt': 0}, {'id': 1, 'name': 'a', 'pt': 1}, + {'id': 2, 'name': None, 'pt': None}, {'id': 4, 'name': 'b', 'pt': 0}, + {'id': 5, 'name': 'a', 'pt': 1}, +] + + +def _create(catalog, *, schema=SCHEMA, partitioned=False, options=None): + settings = { + 'file.format': 'parquet', 'bucket': '-1', 'row-tracking.enabled': 'true', + 'data-evolution.enabled': 'true', 'global-index.enabled': 'true', + 'write.native.enabled': 'true', 'commit.native.enabled': 'false', + 'read.native.enabled': 'false', 'scan.native-plan.enabled': 'true', + 'sorted-index.records-per-range': '2', 'read.batch-size': '1', + 'global-index.search-mode': 'full', + } + settings.update(options or {}) + catalog.create_table('default.indexes', Schema.from_pyarrow_schema( + schema, partition_keys=['pt'] if partitioned else [], options=settings), False) + return catalog.get_table('default.indexes') + + +def _append(table, rows, schema=SCHEMA): + builder = table.copy({'write.native.enabled': 'false'}).new_batch_write_builder() + writer, commit = builder.new_write(), builder.new_commit() + try: + writer.write_arrow(pa.Table.from_pylist(rows, schema=schema)) + commit.commit(writer.prepare_commit()) + finally: + writer.close() + commit.close() + + +def _commit(table, messages): + commit = table.new_batch_write_builder().new_commit() + try: + commit.commit(messages) + finally: + commit.close() + + +def _files(table): + return [entry.index_file for entry in IndexFileHandler(table).scan(None)] + + +def _message_files(messages): + return [entry.index_file for message in messages for entry in message.index_adds] + + +def _path(table, file): + return file.external_path or table.path_factory().global_index_path_factory().to_path(file.file_name) + + +def _ids(table, predicate=None): + builder = table.new_read_builder().with_projection(['id']) + if predicate is not None: + builder.with_filter(predicate) + rows = builder.new_read().to_arrow(builder.new_scan().plan().splits()) + return sorted(rows.column('id').to_pylist()) + + +def _build(table, kind, column='name', **kwargs): + # A native build must never materialize or sort rows in Python. + with patch.object(GlobalIndexBuilder, '_build_sorted_index', side_effect=AssertionError('Python index build ran')): + return GlobalIndexBuilder(table, column, index_type=kind, **kwargs).build() + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +@pytest.mark.parametrize('partitioned', [False, True]) +@pytest.mark.parametrize('native_commit', [False, True]) +def test_build_is_unpublished_and_messages_work_with_either_committer(rest_catalog, kind, partitioned, native_commit): + catalog, _ = rest_catalog + table = _create(catalog, partitioned=partitioned, options={'commit.native.enabled': str(native_commit).lower()}) + _append(table, ROWS) + before = table.snapshot_manager().get_latest_snapshot() + messages = _build(table, kind) + assert messages and all(not message.new_files and message.index_adds for message in messages) + assert table.snapshot_manager().get_latest_snapshot().id == before.id + assert _files(table) == [] + for file in _message_files(messages): + assert file.index_type == kind + assert file.global_index_meta.index_field_id == table.field_dict['name'].id + assert table.file_io.exists(_path(table, file)) + _commit(table, messages) + assert table.snapshot_manager().get_latest_snapshot().id == before.id + 1 + predicate = table.new_read_builder().new_predicate_builder().equal('name', 'a') + with DataEvolutionGlobalIndexScanner.create(table, predicate=predicate) as scanner: + result = scanner.scan(predicate) + assert result is not None and result.is_exact() + assert result.results().cardinality() == 2 + assert _ids(table, predicate) == [1, 5] + nulls = table.new_read_builder().new_predicate_builder().is_null('name') + assert _ids(table, nulls) == [2] + after = table.snapshot_manager().get_latest_snapshot().id + assert _build(table, kind) == [] + assert table.create_global_index('name', index_type=kind) == 0 + assert table.snapshot_manager().get_latest_snapshot().id == after + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +@pytest.mark.parametrize('selection', ['spec', 'predicate', 'null', 'union']) +def test_partition_selection_and_incremental_coverage(rest_catalog, kind, selection): + catalog, _ = rest_catalog + table = _create(catalog, partitioned=True) + _append(table, ROWS) + predicates = table.new_read_builder().new_predicate_builder() + kwargs = ({'partition_filter': predicates.equal('pt', 1)} if selection == 'predicate' else + {'partitions': {'pt': None}} if selection == 'null' else + {'partitions': [{'pt': 0}, {'pt': 1}]} if selection == 'union' else {'partitions': {'pt': 1}}) + first = _build(table, kind, column='id', **kwargs) + expected_partitions = {(None,)} if selection == 'null' else {(0,), (1,)} if selection == 'union' else {(1,)} + assert {message.partition for message in first} == expected_partitions + _commit(table, first) + assert _build(table, kind, column='id', **kwargs) == [] + _append(table, [{'id': 6, 'name': 'a', 'pt': 1}, {'id': 7, 'name': 'z', 'pt': 0}]) + added = _build(table, kind, column='id', **kwargs) + assert sum(file.row_count for file in _message_files(added)) == ( + 0 if selection == 'null' else 2 if selection == 'union' else 1) + if added: + _commit(table, added) + remaining = _build(table, kind, column='id') + assert sum(file.row_count for file in _message_files(remaining)) == ( + 6 if selection == 'null' else 1 if selection == 'union' else 4) + _commit(table, remaining) + assert _build(table, kind, column='id') == [] + assert _ids(table, predicates.greater_or_equal('id', 5)) == [5, 6, 7] + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +def test_external_paths_remain_readable_and_explicit_abort_removes_only_private_files(rest_catalog, tmp_path, kind): + catalog, _ = rest_catalog + root = (tmp_path / 'external-indexes').as_uri() + table = _create(catalog, options={'global-index.external-path': root}) + _append(table, ROWS) + messages = _build(table, kind) + private = _message_files(messages) + assert all(file.external_path.startswith(root + '/') for file in private) + assert all(table.file_io.exists(file.external_path) for file in private) + commit = table.new_batch_write_builder().new_commit() + try: + commit.abort(messages) + finally: + commit.close() + assert all(not table.file_io.exists(file.external_path) for file in private) + assert table.snapshot_manager().get_latest_snapshot().id == 1 + replacement = _build(table, kind) + _commit(table, replacement) + assert {file.file_name for file in private}.isdisjoint(file.file_name for file in _files(table)) + table = table.copy({'global-index.external-path': (tmp_path / 'changed-root').as_uri()}) + assert _ids(table, table.new_read_builder().new_predicate_builder().equal('name', 'a')) == [1, 5] + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +@pytest.mark.parametrize('response', ['first', 'empty']) +def test_rest_snapshot_response_controls_build_instead_of_newer_disk_snapshot(rest_catalog, kind, response): + catalog, server = rest_catalog + table = _create(catalog) + _append(table, ROWS[:2]) + first = table.snapshot_manager().get_latest_snapshot() + _append(table, ROWS[2:]) + reply = (GetTableSnapshotResponse(TableSnapshot(first, 1, 0, 1, first.time_millis)) if response == 'first' + else GetTableSnapshotResponse()) + with patch.object(server, '_table_snapshot_handle', return_value=server._mock_response(reply, 200)): + messages = _build(table, kind) + count = sum(file.row_count for file in _message_files(messages)) + assert count == (2 if response == 'first' else 0) + assert table.snapshot_manager().get_latest_snapshot().id == 2 + + +@pytest.mark.parametrize('code', [403, 500, 503]) +def test_rest_snapshot_errors_propagate_without_starting_python_build(rest_catalog, code): + catalog, server = rest_catalog + table = _create(catalog) + _append(table, ROWS) + reply = ErrorResponse('TABLE', 'indexes', 'snapshot unavailable', code) + with patch.object(server, '_table_snapshot_handle', return_value=server._mock_response(reply, code)): + with pytest.raises(Exception, match='snapshot unavailable|permission'): + _build(table, 'btree') + assert table.snapshot_manager().get_latest_snapshot().id == 1 + assert _files(table) == [] + index_directory = table.path_factory().index_path() + assert not table.file_io.exists(index_directory) or not table.file_io.list_status(index_directory) + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +@pytest.mark.parametrize('column, data_type, values, literal', [ + ('value', pa.int64(), [-2, 3, None, -2], -2), + ('value', pa.bool_(), [True, False, None, True], True), + ('value', pa.date32(), [date(2025, 1, 1), date(2025, 1, 2), None, date(2025, 1, 1)], date(2025, 1, 1)), + ('value', pa.decimal128(12, 2), [Decimal('1.25'), Decimal('2.50'), None, Decimal('1.25')], Decimal('1.25')), + ('value', pa.time32('ms'), [time(1, 2, 3), time(2, 3, 4), None, time(1, 2, 3)], time(1, 2, 3)), + ('value', pa.timestamp('us'), [datetime(2025, 1, 1), datetime(2025, 1, 2), None, datetime(2025, 1, 1)], + datetime(2025, 1, 1)), + ('value', pa.timestamp('us', tz='UTC'), + [datetime(2025, 1, 1, tzinfo=timezone.utc), datetime(2025, 1, 2, tzinfo=timezone.utc), None, + datetime(2025, 1, 1, tzinfo=timezone.utc)], datetime(2025, 1, 1, tzinfo=timezone.utc)), +]) +def test_typed_keys_match_python_scalar_readers(rest_catalog, kind, column, data_type, values, literal): + catalog, _ = rest_catalog + schema = pa.schema([('id', pa.int32()), (column, data_type)]) + table = _create(catalog, schema=schema) + _append(table, [{'id': i, column: value} for i, value in enumerate(values)], schema) + _commit(table, _build(table, kind, column)) + predicate = table.new_read_builder().new_predicate_builder().equal(column, literal) + with DataEvolutionGlobalIndexScanner.create(table, predicate=predicate) as scanner: + result = scanner.scan(predicate) + assert list(result.results()) == [0, 3] + assert _ids(table, predicate) == [0, 3] + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +def test_binding_partition_conjunction_and_rejection_are_core_semantics(rest_catalog, kind): + catalog, _ = rest_catalog + table = _create(catalog, partitioned=True) + _append(table, ROWS) + predicates = table.new_read_builder().new_predicate_builder() + native = _native_table(table).new_sorted_global_index_build_builder().with_index_column('id').with_index_type(kind) + # Python partition predicates carry partition-row indices; bindings resolve by field name. + native.with_partition_filter(_predicate_to_native(predicates.equal('pt', 0))) + native.with_partition_filter(_predicate_to_native(predicates.equal('pt', 1))) + assert native.build() == [] + assert native.execute() == 0 + with pytest.raises(Exception, match='only partition keys'): + native.with_partition_filter(_predicate_to_native(predicates.equal('id', 1))) + assert table.snapshot_manager().get_latest_snapshot().id == 1 + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +def test_empty_table_build_is_a_noop(rest_catalog, kind): + catalog, _ = rest_catalog + table = _create(catalog) + assert _build(table, kind) == [] + assert table.create_global_index('name', index_type=kind) == 0 + assert _native_table(table).new_sorted_global_index_build_builder().with_index_column('name').execute() == 0 + assert table.snapshot_manager().get_latest_snapshot() is None + + +def test_native_binding_can_build_and_execute_separate_indexes(rest_catalog): + catalog, _ = rest_catalog + table = _create(catalog) + _append(table, ROWS) + native = _native_table(table) + builder = native.new_sorted_global_index_build_builder().with_index_column('name') + prepared = builder.build() + assert all(message.serialize() for message in prepared) + _commit(table, from_native_commit_messages(table, prepared)) + assert builder.execute() == 0 + builder.with_index_type('bitmap').with_options({'sorted-index.records-per-range': '100'}) + assert builder.execute() == 1 + assert {file.index_type for file in _files(table)} == {'btree', 'bitmap'} + assert table.snapshot_manager().get_latest_snapshot().id == 3 + + +def test_both_index_types_can_share_one_explicit_commit(rest_catalog): + catalog, _ = rest_catalog + table = _create(catalog) + _append(table, ROWS) + messages = _build(table, 'btree') + _build(table, 'bitmap') + assert table.snapshot_manager().get_latest_snapshot().id == 1 + _commit(table, messages) + assert table.snapshot_manager().get_latest_snapshot().id == 2 + assert {file.index_type for file in _files(table)} == {'btree', 'bitmap'} + assert _ids(table, table.new_read_builder().new_predicate_builder().equal('name', 'a')) == [1, 5] + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +def test_commit_response_failure_keeps_published_index_files(rest_catalog, kind): + catalog, _ = rest_catalog + table = _create(catalog) + _append(table, ROWS) + original = BatchTableCommit.commit + published = [] + + def publish_then_raise(commit, messages, *args, **kwargs): + published.extend(_message_files(messages)) + original(commit, messages, *args, **kwargs) + raise OSError('commit response lost') + + with patch.object(BatchTableCommit, 'commit', publish_then_raise): + with pytest.raises(OSError, match='commit response lost'): + table.create_global_index('name', index_type=kind) + assert published and all(table.file_io.exists(_path(table, file)) for file in published) + assert table.snapshot_manager().get_latest_snapshot().id == 2 + assert _ids(table, table.new_read_builder().new_predicate_builder().equal('name', 'a')) == [1, 5] + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +def test_invalid_native_options_do_not_fall_back_or_publish_files(rest_catalog, kind): + catalog, _ = rest_catalog + table = _create(catalog) + _append(table, ROWS) + with pytest.raises(Exception, match='greater than 0'): + _build(table, kind, options={'sorted-index.records-per-range': 0}) + assert table.snapshot_manager().get_latest_snapshot().id == 1 + index_directory = table.path_factory().index_path() + assert not table.file_io.exists(index_directory) or not table.file_io.list_status(index_directory) + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +def test_java_records_per_file_option_overrides_table_alias(rest_catalog, kind): + catalog, _ = rest_catalog + table = _create(catalog) + _append(table, ROWS) + messages = _build(table, kind, options={'sorted-index.records-per-file': 100}) + files = _message_files(messages) + assert len(files) == 1 + assert files[0].row_count == len(ROWS) + _commit(table, messages) + assert _ids(table, table.new_read_builder().new_predicate_builder().equal('name', 'a')) == [1, 5] + + +@pytest.mark.parametrize('options', [{'write.native.enabled': 'false'}, {'data-evolution.enabled': 'false'}]) +def test_ineligible_tables_choose_python_before_building(rest_catalog, options): + catalog, _ = rest_catalog + table = _create(catalog, options=options) + _append(table, ROWS) + original = GlobalIndexBuilder._build_sorted_index + with patch.object(GlobalIndexBuilder, '_build_sorted_index', autospec=True, side_effect=original) as build: + assert table.create_global_index('name') > 0 + build.assert_called_once() + assert _ids(table, table.new_read_builder().new_predicate_builder().equal('name', 'a')) == [1, 5] + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +@pytest.mark.parametrize('native_commit', [False, True]) +@pytest.mark.parametrize('same_commit', [False, True]) +def test_target_column_changes_reject_staged_indexes_without_deleting_files( + rest_catalog, kind, native_commit, same_commit): + catalog, _ = rest_catalog + table = _create(catalog, options={'commit.native.enabled': str(native_commit).lower()}) + _append(table, ROWS) + messages = _build(table, kind) + private = _message_files(messages) + updater = table.new_batch_write_builder().new_update() + updated = updater.update_by_arrow_with_row_id(pa.table({'_ROW_ID': [0], 'name': ['changed']})) + if not same_commit: + _commit(table, updated) + before = table.snapshot_manager().get_latest_snapshot().id + with pytest.raises(Exception, match='Global index source conflict'): + _commit(table, messages + updated if same_commit else messages) + assert table.snapshot_manager().get_latest_snapshot().id == before + assert all(table.file_io.exists(_path(table, file)) for file in private) + if not same_commit: + assert _ids(table, table.new_read_builder().new_predicate_builder().equal('name', 'changed')) == [3] + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +@pytest.mark.parametrize('native_commit', [False, True]) +def test_append_and_unrelated_column_changes_do_not_invalidate_staged_index(rest_catalog, kind, native_commit): + catalog, _ = rest_catalog + table = _create(catalog, options={'commit.native.enabled': str(native_commit).lower()}) + _append(table, ROWS) + messages = _build(table, kind) + assert all(struct.unpack('>iiq', file.global_index_meta.source_meta) == (0x44454958, 1, 1) + for file in _message_files(messages)) + updater = table.new_batch_write_builder().new_update() + _commit(table, updater.update_by_arrow_with_row_id(pa.table({'_ROW_ID': [0], 'id': pa.array([10], pa.int32())}))) + _append(table, [{'id': 6, 'name': 'a', 'pt': 0}]) + _commit(table, messages) + assert _ids(table, table.new_read_builder().new_predicate_builder().equal('name', 'a')) == [1, 5, 6] + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +@pytest.mark.parametrize('method, literals, partitions', [ + ('in', [], set()), ('notIn', [], {(0,), (1,)}), + ('in', [None, 0], {(0,)}), ('notIn', [None, 0], set()), +]) +def test_partition_set_predicates_follow_java_null_semantics(rest_catalog, kind, method, literals, partitions): + catalog, _ = rest_catalog + table = _create(catalog, partitioned=True) + _append(table, ROWS) + from pypaimon.common.predicate import Predicate + predicate = Predicate(method, 2, 'pt', literals) + messages = _build(table, kind, partition_filter=predicate) + assert {message.partition for message in messages} == partitions + + +@pytest.mark.parametrize('restricted', [False, True]) +def test_query_authorization_remains_on_python_build_path(rest_catalog, restricted): + from pypaimon.catalog.rest.rest_catalog import RESTCatalog + from pypaimon.catalog.table_query_auth import TableQueryAuthResult + from pypaimon.catalog.catalog_exception import TableNoPermissionException + + catalog, _ = rest_catalog + table = _create(catalog) + _append(table, ROWS) + table = table.copy({'query-auth.enabled': 'true'}) + auth = TableQueryAuthResult(None, {'name': '{"name":"NULL"}'} if restricted else None) + original = GlobalIndexBuilder._build_sorted_index + with patch.object(RESTCatalog, 'auth_table_query', return_value=auth), \ + patch('pypaimon.globalindex.native_index_build._native_table', + side_effect=AssertionError('native auth build')), \ + patch.object(GlobalIndexBuilder, '_build_sorted_index', autospec=True, side_effect=original) as build: + if restricted: + with pytest.raises(TableNoPermissionException): + GlobalIndexBuilder(table, 'name').build() + build.assert_not_called() + else: + assert GlobalIndexBuilder(table, 'name').build() + build.assert_called_once() + + +@pytest.mark.parametrize('kind', ['btree', 'bitmap']) +@pytest.mark.parametrize('native_commit', [False, True]) +@pytest.mark.parametrize('remove_all', [False, True]) +def test_deleted_column_deltas_and_removed_source_ranges_reject_staged_index( + rest_catalog, kind, native_commit, remove_all): + from pypaimon.write.commit_message import CommitMessage + + catalog, _ = rest_catalog + table = _create(catalog, options={'commit.native.enabled': str(native_commit).lower()}) + _append(table, ROWS) + updates = table.new_batch_write_builder().new_update().update_by_arrow_with_row_id( + pa.table({'_ROW_ID': [0], 'name': ['changed']})) + _commit(table, updates) + messages = _build(table, kind) + if remove_all: + commit = table.new_batch_write_builder().new_commit() + try: + commit.truncate_table() + finally: + commit.close() + else: + deletions = [CommitMessage(partition=message.partition, bucket=message.bucket, new_files=[], + deleted_files=message.new_files) for message in updates] + _commit(table, deletions) + with pytest.raises(Exception, match='Global index (source|row ID existence) conflict'): + _commit(table, messages) + assert all(table.file_io.exists(_path(table, file)) for file in _message_files(messages)) + + +@pytest.mark.parametrize('native_commit', [False, True]) +@pytest.mark.parametrize('source', [b'DEIX', struct.pack('>iiq', 0x44454958, 1, 99)]) +def test_malformed_or_future_index_source_is_rejected(rest_catalog, native_commit, source): + catalog, _ = rest_catalog + table = _create(catalog, options={'commit.native.enabled': str(native_commit).lower()}) + _append(table, ROWS) + messages = _build(table, 'btree') + for file in _message_files(messages): + file.global_index_meta.source_meta = source + with pytest.raises(Exception, match='source metadata|source snapshot'): + _commit(table, messages) + assert table.snapshot_manager().get_latest_snapshot().id == 1 + assert all(table.file_io.exists(_path(table, file)) for file in _message_files(messages)) diff --git a/paimon-python/pypaimon/write/commit/conflict_detection.py b/paimon-python/pypaimon/write/commit/conflict_detection.py index 19d831cf60e6..179e21ddf7aa 100644 --- a/paimon-python/pypaimon/write/commit/conflict_detection.py +++ b/paimon-python/pypaimon/write/commit/conflict_detection.py @@ -287,6 +287,12 @@ def check_conflicts( if conflict is not None: return conflict + from pypaimon.write.commit.global_index_source_check import check_global_index_sources + conflict = check_global_index_sources( + self, latest_snapshot, merged_entries, delta_entries, delta_index_entries) + if conflict is not None: + return conflict + return self.check_row_id_from_snapshot(latest_snapshot, delta_entries) @staticmethod diff --git a/paimon-python/pypaimon/write/commit/global_index_source_check.py b/paimon-python/pypaimon/write/commit/global_index_source_check.py new file mode 100644 index 000000000000..c4d9bb52ea54 --- /dev/null +++ b/paimon-python/pypaimon/write/commit/global_index_source_check.py @@ -0,0 +1,82 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Revalidate unpublished indexes against the data scanned to build them.""" + +from pypaimon.index.data_evolution_index_source_meta import DataEvolutionIndexSourceMeta +from pypaimon.utils.range import Range + + +def check_global_index_sources(detection, latest_snapshot, merged_entries, delta_entries, index_entries): + if not detection.data_evolution_enabled: + return None + sources = [] + for entry in detection.global_index_file_additions(index_entries): + meta = entry.index_file.global_index_meta + if not DataEvolutionIndexSourceMeta.is_data_evolution_meta(meta.source_meta): + continue + source_id = DataEvolutionIndexSourceMeta.deserialize(meta.source_meta).scan_snapshot_id + if latest_snapshot is None or source_id > latest_snapshot.id: + return RuntimeError('Global index source conflict: source snapshot is not available.') + if detection.snapshot_manager.get_snapshot_by_id(source_id) is None: + return RuntimeError('Global index source conflict: source snapshot {} is missing.'.format(source_id)) + sources.append((source_id, entry)) + if not sources: + return None + conflict = detection.check_global_index_row_id_existence(merged_entries, [entry for _, entry in sources]) + if conflict is not None: + return conflict + conflict = _check_column_changes(detection, sources, delta_entries) + if conflict is not None: + return conflict + earliest = min(source_id for source_id, _ in sources) + for snapshot_id in range(earliest + 1, latest_snapshot.id + 1): + snapshot = detection.snapshot_manager.get_snapshot_by_id(snapshot_id) + if snapshot is None: + return RuntimeError('Global index source conflict: snapshot {} is missing.'.format(snapshot_id)) + if snapshot.commit_kind == 'COMPACT': + continue + changes = detection.commit_scanner.read_incremental_raw_entries_from_changed_partitions( + snapshot, [], index_entries=[entry for _, entry in sources]) + conflict = _check_column_changes(detection, sources, changes, snapshot_id) + if conflict is not None: + return conflict + return None + + +def _check_column_changes(detection, sources, changes, snapshot_id=None): + from pypaimon.write.commit.conflict_detection import RowIdColumnConflictChecker + + field_checker = RowIdColumnConflictChecker([], detection.table.schema_manager) + for change in changes: + row_range = change.file.row_id_range() + if row_range is None: + continue + for source_id, index in sources: + if snapshot_id is not None and snapshot_id <= source_id: + continue + if tuple(index.partition.values) != tuple(change.partition.values) or index.bucket != change.bucket: + continue + meta = index.index_file.global_index_meta + if not row_range.overlaps(Range(meta.row_range_start, meta.row_range_end)): + continue + fields = {meta.index_field_id} + fields.update(meta.extra_field_ids or []) + if field_checker._contains_any_write_field(fields, change.file): + return RuntimeError( + "Global index source conflict: indexed values changed after building " + "index file '{}'.".format(index.index_file.file_name)) + return None From 51f2c37659f14ac8c03c97049ec0557ec70b73b8 Mon Sep 17 00:00:00 2001 From: JingsongLi Date: Sat, 10 Oct 2026 14:11:28 +0800 Subject: [PATCH 2/2] fix(python): repair native vector CI routing and regression coverage --- paimon-python/pypaimon/read/table_read.py | 3 +- .../pypaimon/tests/act_runner_test.py | 28 ++++----- .../tests/batch_vector_lookup_test.py | 17 +++++- .../tests/data_evolution_row_rolling_test.py | 8 ++- .../pypaimon/tests/native_read_test.py | 9 +-- .../pypaimon/tests/native_write_test.py | 58 +++++++++++++++++++ 6 files changed, 96 insertions(+), 27 deletions(-) diff --git a/paimon-python/pypaimon/read/table_read.py b/paimon-python/pypaimon/read/table_read.py index e0e683d973ee..b0b103be4a80 100644 --- a/paimon-python/pypaimon/read/table_read.py +++ b/paimon-python/pypaimon/read/table_read.py @@ -791,8 +791,7 @@ def _read_native_split_group( def _native_split_files_supported(split): for data_file in split.files: file_name = data_file.file_name.lower() - if ('.vector.' not in file_name - and not file_name.endswith(_NATIVE_READ_FILE_SUFFIXES) + if (not file_name.endswith(_NATIVE_READ_FILE_SUFFIXES) and not file_name.endswith(_NATIVE_BLOB_FILE_SUFFIXES)): return False return True diff --git a/paimon-python/pypaimon/tests/act_runner_test.py b/paimon-python/pypaimon/tests/act_runner_test.py index 6e85e516908f..e0f3a5c5b57f 100644 --- a/paimon-python/pypaimon/tests/act_runner_test.py +++ b/paimon-python/pypaimon/tests/act_runner_test.py @@ -258,6 +258,15 @@ def benchmark_input(tmp_path, monkeypatch): return root, warehouse +def _frame_for_append(frames, predicate): + # A full-row append must retain all required Blob fields, including depth + # images outside the RGB-only ACT input projection. + scalar, blobs = frames.scan().where(predicate).read_blobs() + row = scalar.to_pylist()[0] + row.update({name: values[0] for name, values in blobs.items()}) + return row + + class _Policy(torch.nn.Module): def __init__(self): @@ -630,15 +639,10 @@ def test_paimon_windows_are_lazy_snapshot_pinned_and_vortex_independent( for name in IMAGE_COLUMNS } == {name: 1 for name in IMAGE_COLUMNS} - scalar, blobs = frames.scan().where( - "episode_id = 'train-a' AND frame_index = 5" - ).read_blobs(IMAGE_COLUMNS) - appended = scalar.to_pylist()[0] + appended = _frame_for_append(frames, "episode_id = 'train-a' AND frame_index = 5") appended["frame_index"] = 6 appended["index"] = max( row["index"] for row in frames.scan().select(["index"]).to_list()) + 1 - for name in IMAGE_COLUMNS: - appended[name] = blobs[name][0] frames.add([appended]) assert train.snapshot_id == snapshot_id @@ -749,12 +753,8 @@ def test_preparation_uses_published_group_after_new_episode_commit( changed["split"] = "test" episodes.add([changed]) frames = connection.get_table(agilex.FRAMES_TABLE) - scalar, blobs = frames.scan().where( - "episode_id = 'train-a' AND frame_index = 0").read_blobs(IMAGE_COLUMNS) - changed_frame = scalar.to_pylist()[0] + changed_frame = _frame_for_append(frames, "episode_id = 'train-a' AND frame_index = 0") changed_frame["action"] = [0.0] * 14 - for name in IMAGE_COLUMNS: - changed_frame[name] = blobs[name][0] frames.add([changed_frame]) assert latest_snapshot_id(frames) != before["paimon"]["frames_snapshot_id"] @@ -810,12 +810,8 @@ def test_dense_comparison_rejects_selected_episode_with_invalid_frames( options={"warehouse": str(warehouse)}) frames = connection.get_table(agilex.FRAMES_TABLE) indices = episode_indices(connection, "act-test@1") - scalar, blobs = frames.scan().where( - "episode_id = 'val-a' AND frame_index = 0").read_blobs(IMAGE_COLUMNS) - invalid = scalar.to_pylist()[0] + invalid = _frame_for_append(frames, "episode_id = 'val-a' AND frame_index = 0") invalid["quality_status"] = 2 - for name in IMAGE_COLUMNS: - invalid[name] = blobs[name][0] frames.add([invalid]) frames.raw_table.create_tag("invalid-frames") diff --git a/paimon-python/pypaimon/tests/batch_vector_lookup_test.py b/paimon-python/pypaimon/tests/batch_vector_lookup_test.py index b849fa1a64f4..229c9cf3a5eb 100644 --- a/paimon-python/pypaimon/tests/batch_vector_lookup_test.py +++ b/paimon-python/pypaimon/tests/batch_vector_lookup_test.py @@ -77,7 +77,9 @@ def read(table_read, splits, *args, **kwargs): assert len(actual) == len(expected) for left, right in zip(expected, actual): assert left.schema == right.schema - assert left.to_pylist() == right.to_pylist() + # Unordered lookups may visit partition files in different orders. + # Keep duplicates and compare all projected values, including Blobs. + assert sorted(map(repr, left.to_pylist())) == sorted(map(repr, right.to_pylist())) if projection == [] and not with_row_id: assert right.num_columns == 0 assert right.num_rows == 3 @@ -158,11 +160,20 @@ def test_single_query_keeps_single_lookup(docs): assert result[0].equals(expected) -def test_unsupported_nested_row_projection_still_raises(docs): +@pytest.mark.python_plan +@pytest.mark.python_read +def test_python_nested_row_projection_still_raises(docs): with pytest.raises(NotImplementedError, match="ROW nested-field projection"): docs.search_vectors([[0.0, 0.0], [1.0, 0.0]]).select(["info.label"]).to_arrow() +@pytest.mark.native_plan +def test_native_nested_row_projection(docs): + docs.raw_table = docs.raw_table.copy({'read.native.enabled': 'true', 'scan.native-plan.enabled': 'true'}) + results = docs.search_vectors([[0.0, 0.0], [1.0, 0.0]]).select(["info.label"]).limit(1).to_arrow() + assert [result.to_pylist() for result in results] == [[{'info_label': '0'}], [{'info_label': '1'}]] + + def test_native_index_batch_uses_shared_lookup(docs): pytest.importorskip("paimon_vindex") docs.raw_table.copy({"deletion-vectors.enabled": "false"}).create_global_index( @@ -180,5 +191,5 @@ def read(execution, result): with patch.object(ScanQuery, "_read_global_index_result", read): actual = query.to_list() - assert actual == expected + assert [sorted(map(repr, rows)) for rows in actual] == [sorted(map(repr, rows)) for rows in expected] assert calls == [6] diff --git a/paimon-python/pypaimon/tests/data_evolution_row_rolling_test.py b/paimon-python/pypaimon/tests/data_evolution_row_rolling_test.py index e37dac3d00c3..e38d05cffa8a 100644 --- a/paimon-python/pypaimon/tests/data_evolution_row_rolling_test.py +++ b/paimon-python/pypaimon/tests/data_evolution_row_rolling_test.py @@ -646,7 +646,9 @@ def test_vector_writer_supports_target_file_row_num(self): 'vector.file.format': 'parquet', }) - files = self._write_files(table, self._vector_rows(7)) + # Native writers check the row limit at Arrow batch boundaries. + data = pa.Table.from_batches(self._vector_rows(7).to_batches(max_chunksize=3)) + files = self._write_files(table, data) data_rows = sorted( f.row_count for f in files @@ -667,7 +669,9 @@ def test_dedicated_writer_rolls_blob_and_vector_together(self): 'vector.file.format': 'parquet', }) - files = self._write_files(table, self._blob_vector_rows(7)) + # Keep all dedicated writers aligned while rolling between batches. + data = pa.Table.from_batches(self._blob_vector_rows(7).to_batches(max_chunksize=3)) + files = self._write_files(table, data) data_rows = sorted( f.row_count for f in files diff --git a/paimon-python/pypaimon/tests/native_read_test.py b/paimon-python/pypaimon/tests/native_read_test.py index 3e117cb0a0c0..094eb89aaf56 100644 --- a/paimon-python/pypaimon/tests/native_read_test.py +++ b/paimon-python/pypaimon/tests/native_read_test.py @@ -951,10 +951,11 @@ def test_native_avro_read_uses_native_path(): native.assert_called_once() -def test_native_read_falls_back_for_unsupported_dedicated_file(): +@pytest.mark.parametrize('file_name', ['camera.unsupported', 'data.vector.lance', 'data.vector.vortex']) +def test_native_read_falls_back_for_unsupported_dedicated_file(file_name): read = _table_read() schema = pa.schema([('id', pa.int32())]) - split = _Split('camera.unsupported') + split = _Split(file_name) split._native_split = object() with patch('pypaimon.read.native_plan.native_read') as native: @@ -963,8 +964,8 @@ def test_native_read_falls_back_for_unsupported_dedicated_file(): native.assert_not_called() -@pytest.mark.parametrize('file_name', ['picture.blob', 'camera.video']) -def test_native_read_supports_blob_and_video_files_and_forwards_parallelism(file_name): +@pytest.mark.parametrize('file_name', ['picture.blob', 'camera.video', 'data.vector.parquet']) +def test_native_read_supports_dedicated_files_and_forwards_parallelism(file_name): read = _table_read() schema = pa.schema([('id', pa.int32())]) split = _Split(file_name) diff --git a/paimon-python/pypaimon/tests/native_write_test.py b/paimon-python/pypaimon/tests/native_write_test.py index dafe004c3763..1d1f0bbc3662 100644 --- a/paimon-python/pypaimon/tests/native_write_test.py +++ b/paimon-python/pypaimon/tests/native_write_test.py @@ -129,6 +129,64 @@ def test_batch_native_write_commits_through_both_committers( assert table.snapshot_manager().get_latest_snapshot().commit_user == builder.commit_user +@requires_native +@pytest.mark.parametrize('evolution', [False, True]) +@pytest.mark.parametrize('streaming', [False, True]) +@pytest.mark.parametrize('element_type', [pa.float32(), pa.float64()]) +@pytest.mark.parametrize('child_name', ['item', 'element', '']) +def test_native_vector_write_preserves_slices_and_nulls( + native_rest_catalog, evolution, streaming, element_type, child_name): + schema = pa.schema([ + ('id', pa.int32()), ('embedding', pa.list_(pa.field(child_name, element_type), 2)), + ]) + options = { + 'file.format': 'parquet', + 'write.native.enabled': 'true', 'commit.native.enabled': 'true', + 'row-tracking.enabled': str(evolution).lower(), + 'data-evolution.enabled': str(evolution).lower(), + } + if evolution: + options['vector.file.format'] = 'parquet' + native_rest_catalog.create_table('default.vectors', Schema.from_pyarrow_schema(schema, options=options), False) + table = native_rest_catalog.get_table('default.vectors') + builder = table.new_stream_write_builder() if streaming else table.new_batch_write_builder() + writer = builder.new_write() + commit = builder.new_commit() + data = pa.Table.from_pylist([ + {'id': 0, 'embedding': [0.0, 1.0]}, {'id': 1, 'embedding': [2.0, 3.0]}, + {'id': 2, 'embedding': None}, {'id': 3, 'embedding': [4.0, None]}, + ], schema=schema).slice(1) + try: + assert isinstance(writer, NativeTableWrite) + assert writer._python_writer is None + with patch.object(NativeTableWrite, '_switch_to_python', side_effect=AssertionError('Vector fallback')): + writer.write_arrow(data) + messages = writer.prepare_commit(7) if streaming else writer.prepare_commit() + commit.commit(messages, 7) if streaming else commit.commit(messages) + finally: + writer.close() + commit.close() + assert _rows(table) == data.to_pylist() + + +@requires_native +@pytest.mark.parametrize('streaming', [False, True]) +def test_unavailable_vector_format_falls_back_before_accepting_data(native_rest_catalog, streaming): + schema = pa.schema([('id', pa.int32()), ('embedding', pa.list_(pa.float32(), 2))]) + native_rest_catalog.create_table('default.vectors', Schema.from_pyarrow_schema(schema, options={ + 'write.native.enabled': 'true', 'row-tracking.enabled': 'true', + 'data-evolution.enabled': 'true', 'vector.file.format': 'lance', + }), False) + table = native_rest_catalog.get_table('default.vectors') + builder = table.new_stream_write_builder() if streaming else table.new_batch_write_builder() + writer = builder.new_write() + try: + assert not isinstance(writer, NativeTableWrite) + assert table.snapshot_manager().get_latest_snapshot() is None + finally: + writer.close() + + @requires_native @pytest.mark.parametrize('directory', [None, 'relative', 'absolute', 'uri']) def test_escaped_partition_file_path_and_abort(tmp_path, native_rest_catalog, directory):