Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion paimon-python/pypaimon/globalindex/create_global_index.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
49 changes: 49 additions & 0 deletions paimon-python/pypaimon/globalindex/native_index_build.py
Original file line number Diff line number Diff line change
@@ -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())
48 changes: 48 additions & 0 deletions paimon-python/pypaimon/index/data_evolution_index_source_meta.py
Original file line number Diff line number Diff line change
@@ -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)
3 changes: 1 addition & 2 deletions paimon-python/pypaimon/read/table_read.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
28 changes: 12 additions & 16 deletions paimon-python/pypaimon/tests/act_runner_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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"]

Expand Down Expand Up @@ -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")

Expand Down
17 changes: 14 additions & 3 deletions paimon-python/pypaimon/tests/batch_vector_lookup_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand All @@ -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]
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
9 changes: 5 additions & 4 deletions paimon-python/pypaimon/tests/native_read_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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)
Expand Down
Loading
Loading