blob: 1387b02536bcd7f6ba227567f4e4fffff9c37362 [file]
# 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.
"""Full-text scan to scan index files."""
from abc import ABC, abstractmethod
from collections import defaultdict
from typing import List
from pypaimon.globalindex.data_evolution_global_index_coverage import DataEvolutionGlobalIndexCoverage
from pypaimon.globalindex.full_text.native_full_text_global_index_reader import (
FULL_TEXT_IDENTIFIER,
)
from pypaimon.table.source.full_text_search_split import (
FullTextSearchSplit,
IndexFullTextSearchSplit,
RawFullTextSearchSplit,
)
from pypaimon.utils.range import Range
class FullTextScanPlan:
"""Plan of full-text scan."""
def __init__(self, splits: List[FullTextSearchSplit]):
self._splits = splits
def splits(self) -> List[FullTextSearchSplit]:
return self._splits
class FullTextScan(ABC):
"""Full-text scan to scan index files."""
@abstractmethod
def scan(self) -> FullTextScanPlan:
pass
class DataEvolutionFullTextScan(FullTextScan):
"""Implementation for FullTextScan."""
def __init__(
self,
table: 'FileStoreTable',
text_columns,
partition_filter=None):
self._table = table
self._text_columns = list(text_columns)
self._partition_filter = partition_filter
def scan(self) -> FullTextScanPlan:
from pypaimon.index.index_file_handler import IndexFileHandler
if not self._text_columns:
return FullTextScanPlan([])
text_column_ids = {field.id for field in self._text_columns}
id_to_column = {field.id: field.name for field in self._text_columns}
from pypaimon.snapshot.time_travel_util import TimeTravelUtil
from pypaimon.common.options.options import Options
snapshot = TimeTravelUtil.try_travel_to_snapshot(
Options(self._table.table_schema.options),
self._table.tag_manager(),
self._table.snapshot_manager(),
)
if snapshot is None:
snapshot = self._table.snapshot_manager().get_latest_snapshot()
index_file_handler = IndexFileHandler(table=self._table)
partition_filter = self._partition_filter
def index_file_filter(entry):
if partition_filter is not None:
if not partition_filter.test(entry.partition):
return False
global_index_meta = entry.index_file.global_index_meta
if global_index_meta is None:
return False
return (
global_index_meta.index_field_id in text_column_ids
and _supports_full_text_search(entry.index_file.index_type)
)
entries = index_file_handler.scan(snapshot, index_file_filter)
all_index_files = [entry.index_file for entry in entries]
# Group full-text index files by column and (rowRangeStart, rowRangeEnd).
by_column_and_range = defaultdict(lambda: defaultdict(list))
for index_file in all_index_files:
meta = index_file.global_index_meta
assert meta is not None
range_key = Range(meta.row_range_start, meta.row_range_end)
column_name = id_to_column[meta.index_field_id]
by_column_and_range[column_name][range_key].append(index_file)
splits = []
for column_name, by_range in by_column_and_range.items():
for range_key, files in by_range.items():
splits.append(
IndexFullTextSearchSplit(
column_name, range_key.from_, range_key.to, files))
if all_index_files:
raw_row_ranges = DataEvolutionGlobalIndexCoverage(
self._table,
snapshot,
partition_filter,
all_index_files,
).unindexed_ranges(
list(text_column_ids),
search_mode=self._table.options.full_text_index_search_mode(),
)
if raw_row_ranges:
splits.append(RawFullTextSearchSplit(raw_row_ranges))
return FullTextScanPlan(splits)
def _supports_full_text_search(index_type):
return index_type == FULL_TEXT_IDENTIFIER