diff --git a/cpp/src/cwrapper/tsfile_cwrapper.cc b/cpp/src/cwrapper/tsfile_cwrapper.cc index ff70990ca..c0fd7f842 100644 --- a/cpp/src/cwrapper/tsfile_cwrapper.cc +++ b/cpp/src/cwrapper/tsfile_cwrapper.cc @@ -567,7 +567,7 @@ ERRNO tsfile_writer_write_arrow(TsFileWriter writer, ArrowArray* array, // Query PreparedSeriesHandle tsfile_reader_prepare_series( - TsFileReader reader, const TsFilePreparedLocator* locator, + TsFileReader reader, const TsFilePreparedLocator* locator, bool trust_index, ERRNO* err_code) { if (err_code == nullptr) { return nullptr; @@ -581,6 +581,7 @@ PreparedSeriesHandle tsfile_reader_prepare_series( generation.file_id = locator->file_id; generation.file_size = locator->file_size; generation.file_fingerprint = locator->file_fingerprint; + generation.trust_index = trust_index; storage::PreparedLocator native_locator; native_locator.locator_id = locator->locator_id; native_locator.layout = locator->layout; @@ -605,7 +606,8 @@ PreparedSeriesHandle tsfile_reader_prepare_series( PreparedSeriesHandle tsfile_reader_prepare_series_with_time_owner( TsFileReader reader, const TsFilePreparedLocator* locator, - PreparedSeriesHandle aligned_time_owner, ERRNO* err_code) { + PreparedSeriesHandle aligned_time_owner, bool trust_index, + ERRNO* err_code) { if (err_code == nullptr) { return nullptr; } @@ -620,6 +622,7 @@ PreparedSeriesHandle tsfile_reader_prepare_series_with_time_owner( generation.file_id = locator->file_id; generation.file_size = locator->file_size; generation.file_fingerprint = locator->file_fingerprint; + generation.trust_index = trust_index; storage::PreparedLocator native_locator; native_locator.locator_id = locator->locator_id; native_locator.layout = locator->layout; diff --git a/cpp/src/cwrapper/tsfile_cwrapper.h b/cpp/src/cwrapper/tsfile_cwrapper.h index 4a87db90c..ee1158f61 100644 --- a/cpp/src/cwrapper/tsfile_cwrapper.h +++ b/cpp/src/cwrapper/tsfile_cwrapper.h @@ -729,15 +729,19 @@ ERRNO tsfile_writer_write_arrow(TsFileWriter writer, ArrowArray* array, /*-------------------TsFile reader query data------------------ */ -/** Deserialize one exact Dataset Index locator into a reusable series. */ +/** Deserialize one exact Dataset Index locator into a reusable series. + * trust_index skips source generation checks; locator bounds remain checked. + */ PreparedSeriesHandle tsfile_reader_prepare_series( - TsFileReader reader, const TsFilePreparedLocator* locator, ERRNO* err_code); + TsFileReader reader, const TsFilePreparedLocator* locator, bool trust_index, + ERRNO* err_code); /** Prepare an aligned value locator by sharing an existing parsed time index. + * trust_index skips source generation checks; locator bounds remain checked. */ PreparedSeriesHandle tsfile_reader_prepare_series_with_time_owner( TsFileReader reader, const TsFilePreparedLocator* locator, - PreparedSeriesHandle aligned_time_owner, ERRNO* err_code); + PreparedSeriesHandle aligned_time_owner, bool trust_index, ERRNO* err_code); /** Release a prepared handle. Existing result sets remain independently owned. */ diff --git a/cpp/src/file/tsfile_io_reader.cc b/cpp/src/file/tsfile_io_reader.cc index 49f5fd9f8..0db0eb228 100644 --- a/cpp/src/file/tsfile_io_reader.cc +++ b/cpp/src/file/tsfile_io_reader.cc @@ -142,13 +142,16 @@ int TsFileIOReader::prepare_series( uint64_t actual_size = 0; uint64_t actual_fingerprint = 0; if (read_file_ == nullptr || generation.file_size == 0 || - read_file_->generation(actual_size, actual_fingerprint) != E_OK || - actual_size != generation.file_size || - (generation.file_fingerprint != 0 && - actual_fingerprint != generation.file_fingerprint) || locator.value_metadata_length == 0 || locator.layout > 1) { return E_INVALID_ARG; } + if (!generation.trust_index && + (read_file_->generation(actual_size, actual_fingerprint) != E_OK || + actual_size != generation.file_size || + (generation.file_fingerprint != 0 && + actual_fingerprint != generation.file_fingerprint))) { + return E_INVALID_ARG; + } std::shared_ptr candidate = std::make_shared(generation, locator); TimeseriesIndex* value_index = nullptr; diff --git a/cpp/src/reader/prepared_series.h b/cpp/src/reader/prepared_series.h index ce647966b..5df4ebeca 100644 --- a/cpp/src/reader/prepared_series.h +++ b/cpp/src/reader/prepared_series.h @@ -34,12 +34,15 @@ struct FileGeneration { uint32_t file_id; uint64_t file_size; uint64_t file_fingerprint; + // A sealed Dataset may skip checking the source file generation. + bool trust_index; FileGeneration() : mapped_index_identity(0), file_id(0), file_size(0), - file_fingerprint(0) {} + file_fingerprint(0), + trust_index(false) {} }; struct PreparedLocator { diff --git a/python/README-zh.md b/python/README-zh.md index 9d322370d..59d2eb0f1 100644 --- a/python/README-zh.md +++ b/python/README-zh.md @@ -68,6 +68,67 @@ mvn -P with-cpp,with-python clean verify python setup.py build_ext --inplace ``` +## Dataset 索引 + +通过 `use_index=True` 启用持久化 Dataset 索引。第一次打开时构建索引,之后 +复用已有索引。`trust_index=True` 是默认值,假设索引对应的 TsFile 不再变更, +因此加载索引和查询时跳过源文件大小、修改时间和文件代次检查。 +索引格式与 locator 边界始终会检查。 + +```python +from tsfile import TsFileDataFrame + +with TsFileDataFrame("dataset/", use_index=True) as dataset: + values = dataset[0][:] + +# 加载索引和获取查询 reader 时检查源文件是否发生变更。 +with TsFileDataFrame("dataset/", use_index=True, trust_index=False) as dataset: + values = dataset[0][:] +``` + +设置 `trust_index=False` 时,打开 dataset 会重建已过期的索引;获取查询 reader +时检测到文件变更则报错。`trust_index` 仅限关键字传入,在默认的 +`use_index=False` 模式下没有作用。 + +## Dataset 读取资源上限 + +`TsFileDataFrame` 在 `use_index=True` 时支持以下仅限关键字的配置参数: + +```python +from tsfile import TsFileDataFrame + +with TsFileDataFrame( + ["part1.tsfile", "part2.tsfile"], + use_index=True, + max_prepared_series=32, + descriptor_cache_size=32, + max_open_files=16, + query_workers=4, +) as frame: + values = frame[0][:] +``` + +| 参数 | 环境变量 | 默认值 | 有效取值 | +|------|----------|--------|----------| +| `max_prepared_series` | `TSFILE_DATAFRAME_MAX_PREPARED_SERIES` | 4096 | 非负整数 | +| `descriptor_cache_size` | `TSFILE_DATAFRAME_DESCRIPTOR_CACHE_SIZE` | 4096 | 非负整数 | +| `max_open_files` | `TSFILE_DATAFRAME_MAX_OPEN_FILES` | 16 | 正整数 | +| `query_workers` | `TSFILE_DATAFRAME_QUERY_WORKERS` | `min(4, os.cpu_count() or 1)` | 正整数 | +| `query_parallel_min_rows` | `TSFILE_DATAFRAME_QUERY_PARALLEL_MIN_ROWS` | 8192 | 正整数 | + +显式参数优先于对应的环境变量。参数为 `None` 时读取环境变量;变量未设置时 +使用内置默认值。配置在构造 frame 时一次性确定,之后修改环境变量不会影响 +已有 frame。子集共享父 frame 的 runtime 和配置。非法值在打开数据文件或索引前 +报错。显式传入这些参数要求 `use_index=True`;默认的无索引模式不读取这些环境变量。 + +预备序列缓存保留已解析的原生序列元数据,按 LRU 淘汰空闲条目并释放原生句柄。 +查询正在使用的条目,以及其他序列仍依赖的共享时间元数据,会保持有效,因而 +并发查询期间条目数可能暂时超过上限,查询结束后再回收。这是每个 runtime 的 +条目数量上限,不是整个进程的内存上限;设置为 `0` 表示使用结束后不保留缓存。 +`descriptor_cache_size` 分别限制名称描述符缓存和序列路由缓存,`0` 禁用二者。 +`max_open_files` 限制打开的 reader 数量,`query_workers=1` 表示串行执行查询组。 +多个 runtime 或工作进程分别维护各自的上限。 + ## 文件级 Properties `TsFileWriter` 和 `TsFileTableWriter` 可以在打开期间写入二进制 property。 diff --git a/python/README.md b/python/README.md index 22e063166..60607fe7d 100644 --- a/python/README.md +++ b/python/README.md @@ -61,6 +61,50 @@ Build by python command: python setup.py build_ext --inplace ``` +## Dataset read resource limits + +`TsFileDataFrame` accepts keyword-only resource options when `use_index=True`: + +```python +from tsfile import TsFileDataFrame + +with TsFileDataFrame( + ["part1.tsfile", "part2.tsfile"], + use_index=True, + max_prepared_series=32, + descriptor_cache_size=32, + max_open_files=16, + query_workers=4, +) as frame: + values = frame[0][:] +``` + +| Parameter | Environment variable | Default | Allowed values | +|-----------|----------------------|---------|----------------| +| `max_prepared_series` | `TSFILE_DATAFRAME_MAX_PREPARED_SERIES` | 4096 | Integer >= 0 | +| `descriptor_cache_size` | `TSFILE_DATAFRAME_DESCRIPTOR_CACHE_SIZE` | 4096 | Integer >= 0 | +| `max_open_files` | `TSFILE_DATAFRAME_MAX_OPEN_FILES` | 16 | Integer >= 1 | +| `query_workers` | `TSFILE_DATAFRAME_QUERY_WORKERS` | `min(4, os.cpu_count() or 1)` | Integer >= 1 | +| `query_parallel_min_rows` | `TSFILE_DATAFRAME_QUERY_PARALLEL_MIN_ROWS` | 8192 | Integer >= 1 | + +An explicit parameter takes precedence over its environment variable. `None` +uses the environment variable, or the built-in default if it is unset. Options +are resolved once when constructing the frame; later environment changes do not +affect it. Subsets share their parent's runtime and limits. Invalid values raise +an error before dataset files or indexes are opened. Explicit resource options +require `use_index=True`; the default non-indexed mode ignores these environment +variables. + +The prepared-series cache retains parsed native series metadata. It evicts the +least recently used idle entries and releases their native handles. Active +queries and dependent series pin shared time metadata, so the entry count can +temporarily exceed the limit until those queries finish. This is an entry-count +limit per runtime, not a process-wide memory limit. A value of `0` disables +retention after use. `descriptor_cache_size` bounds each of the named-descriptor +and series-route caches; `0` disables both. `max_open_files` bounds open readers, +and `query_workers=1` runs query groups serially. Multiple runtimes or worker +processes have separate limits. + ## File-level properties `TsFileWriter` and `TsFileTableWriter` accept binary properties while they are @@ -78,6 +122,29 @@ with TsFileReader("example.tsfile") as reader: Values do not carry a data type; use an explicit portable encoding when storing numbers or structures. +## Dataset indexes + +Enable the persistent dataset index with `use_index=True`. The first open +builds the index; subsequent opens reuse it. By default, `trust_index=True` +assumes the indexed TsFiles will not change and skips source file size, +modification time, and generation checks during index loading and queries. +Index format and locator bounds are always checked. + +```python +from tsfile import TsFileDataFrame + +with TsFileDataFrame("dataset/", use_index=True) as dataset: + series = dataset[0] + +# Check source generations when loading the index and acquiring readers. +with TsFileDataFrame("dataset/", use_index=True, trust_index=False) as dataset: + series = dataset[0] +``` + +With `trust_index=False`, a stale index is rebuilt when the dataset opens; +changes detected while acquiring a query reader raise an error. `trust_index` +is keyword-only and has no effect when `use_index=False`, which remains the default. + ## Local File Read Backend Python readers inherit the process-wide backend setting when they open a file. diff --git a/python/tests/test_dataset_index.py b/python/tests/test_dataset_index.py index 3d768540e..750fdb119 100644 --- a/python/tests/test_dataset_index.py +++ b/python/tests/test_dataset_index.py @@ -640,6 +640,67 @@ def fail_legacy_scan(*_args, **_kwargs): series.close() +@pytest.mark.parametrize("capacity", [0, 1, 2]) +def test_trust_index_defaults_to_skipping_generation_checks( + tmp_path, monkeypatch, capacity +): + source = tmp_path / "part.tsfile" + _write_runtime_file(source, 0) + with TsFileDataFrame(str(source), show_progress=False, use_index=True): + pass + + stat = os.stat(source) + os.utime(source, ns=(stat.st_atime_ns, stat.st_mtime_ns + 1_000_000)) + + def fail_check(*_args, **_kwargs): + raise AssertionError("trusted indexes must not check source generations") + + monkeypatch.setattr(index_module, "file_fingerprint", fail_check) + monkeypatch.setattr(runtime_module, "file_fingerprint", fail_check) + monkeypatch.setattr( + runtime_module._ReaderSession, "_validate_generation", fail_check + ) + monkeypatch.setattr("tsfile.dataset.reader.TsFileSeriesReader", fail_check) + with TsFileDataFrame( + str(source), + show_progress=False, + use_index=True, + max_prepared_series=capacity, + descriptor_cache_size=capacity, + ) as dataframe: + with dataframe[:1] as subset: + assert subset._trust_index is True + np.testing.assert_array_equal(subset[0][:], np.array([0.0, 1.0])) + np.testing.assert_array_equal(dataframe[0][:], np.array([0.0, 1.0])) + assert dataframe._runtime.prepared.size <= capacity + + +def test_trust_index_false_rebuilds_a_stale_index(tmp_path): + source = tmp_path / "part.tsfile" + _write_runtime_file(source, 0) + with TsFileDataFrame(str(source), show_progress=False, use_index=True) as dataframe: + index_path = dataframe._runtime.index.path + fingerprint = dataframe._runtime.index.record(TSFILE_RECORD, 0)[3] + + stat = os.stat(source) + os.utime(source, ns=(stat.st_atime_ns, stat.st_mtime_ns + 1_000_000)) + assert index_module.index_matches_paths(index_path, [str(source)]) + assert not index_module.index_matches_paths( + index_path, [str(source)], trust_index=False + ) + with TsFileDataFrame( + str(source), show_progress=False, use_index=True, trust_index=False + ) as dataframe: + assert dataframe._runtime.index.record(TSFILE_RECORD, 0)[3] != fingerprint + np.testing.assert_array_equal(dataframe[0][:], np.array([0.0, 1.0])) + + +@pytest.mark.parametrize("trust_index", [None, 1, "true"]) +def test_trust_index_requires_a_bool(tmp_path, trust_index): + with pytest.raises(TypeError, match="trust_index must be a bool"): + TsFileDataFrame(str(tmp_path / "missing.tsfile"), trust_index=trust_index) + + def test_dataframe_does_not_use_or_create_index_by_default(tmp_path): source = tmp_path / "part.tsfile" _write_runtime_file(source, 0) @@ -797,9 +858,9 @@ def count_find_device(*args, **kwargs): def test_runtime_descriptor_cache_evicts_least_recent_name(tmp_path, monkeypatch): source = tmp_path / "devices.tsfile" _write_runtime_devices_file(source) - monkeypatch.setattr(runtime_module, "_SERIES_DESCRIPTOR_CACHE_SIZE", 2) - - with TsFileDataFrame(str(source), show_progress=False, use_index=True) as dataframe: + with TsFileDataFrame( + str(source), show_progress=False, use_index=True, descriptor_cache_size=2 + ) as dataframe: names = [str(name) for name in dataframe.list_timeseries()] find_device_calls = 0 original_find_device = dataframe._runtime.index.find_device_id @@ -925,8 +986,9 @@ def test_prepared_query_reads_nullable_offset_window_in_arrow_batches(tmp_path): runtime = dataframe._runtime series = runtime.index.record(LOGICAL_SERIES, 0) span = runtime.index.record(SERIES_FILE_SPAN, series[2]) - with runtime.readers.acquire(0) as reader: - prepared = runtime.prepared.get(0, span[2], reader) + with runtime.readers.acquire(0) as reader, runtime.prepared.acquire( + 0, span[2], reader + ) as prepared: with reader.query_prepared(prepared, offset=1, limit=7) as result: batches = [] while True: @@ -951,7 +1013,8 @@ def test_prepared_query_reads_nullable_offset_window_in_arrow_batches(tmp_path): ) -def test_prepared_locator_rejects_stale_generation_and_bad_range(tmp_path): +@pytest.mark.parametrize("trust_index", [True, False]) +def test_prepared_locator_generation_checks_and_bounds(tmp_path, trust_index): source = tmp_path / "part.tsfile" _write_runtime_file(source, 0) with TsFileDataFrame(str(source), show_progress=False, use_index=True) as dataframe: @@ -962,27 +1025,38 @@ def test_prepared_locator_rejects_stale_generation_and_bad_range(tmp_path): with runtime.readers.acquire(0) as reader: stale = list(locator) stale[3] ^= 1 - with pytest.raises(Exception, match="prepare Dataset Index locator"): - reader.prepare_series(stale) + if trust_index: + reader.prepare_series(stale, trust_index=True).close() + else: + with pytest.raises(Exception, match="prepare Dataset Index locator"): + reader.prepare_series(stale) out_of_range = list(locator) out_of_range[7] = os.path.getsize(source) + 1 with pytest.raises(Exception, match="prepare Dataset Index locator"): - reader.prepare_series(out_of_range) + reader.prepare_series(out_of_range, trust_index=trust_index) -def test_reader_session_revalidates_generation_when_reused(tmp_path): +@pytest.mark.parametrize("trust_index", [True, False]) +def test_reader_session_generation_checks_follow_trust_index(tmp_path, trust_index): source = tmp_path / "part.tsfile" _write_runtime_file(source, 0) - with TsFileDataFrame(str(source), show_progress=False, use_index=True) as dataframe: + with TsFileDataFrame( + str(source), show_progress=False, use_index=True, trust_index=trust_index + ) as dataframe: pool = dataframe._runtime.readers with pool.acquire(0): pass stat = os.stat(source) os.utime(source, ns=(stat.st_atime_ns, stat.st_mtime_ns + 1_000_000)) - with pytest.raises(RuntimeError, match="generation changed"): + if trust_index: with pool.acquire(0): pass + np.testing.assert_array_equal(dataframe[0][:], np.array([0.0, 1.0])) + else: + with pytest.raises(RuntimeError, match="generation changed"): + with pool.acquire(0): + pass def test_runtime_lease_close_waits_for_query_lease(tmp_path): diff --git a/python/tests/test_dataset_read_options.py b/python/tests/test_dataset_read_options.py new file mode 100644 index 000000000..8872f8cc8 --- /dev/null +++ b/python/tests/test_dataset_read_options.py @@ -0,0 +1,524 @@ +# 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. + +from collections import OrderedDict +from concurrent.futures import ThreadPoolExecutor +from contextlib import ExitStack +import threading +from types import SimpleNamespace + +import numpy as np +import pandas as pd +import pytest + +from tsfile import ( + ColumnCategory, + ColumnSchema, + TableSchema, + TSDataType, + TsFileDataFrame, + TsFileTableWriter, +) +from tsfile.dataset.runtime import PreparedSeriesCache + +_OPTIONS = ( + "max_prepared_series", + "descriptor_cache_size", + "max_open_files", + "query_workers", + "query_parallel_min_rows", +) + + +@pytest.fixture +def dataset_file(tmp_path): + path = tmp_path / "devices.tsfile" + schema = TableSchema( + "weather", + [ + ColumnSchema("device", TSDataType.STRING, ColumnCategory.TAG), + ColumnSchema("value", TSDataType.DOUBLE, ColumnCategory.FIELD), + ColumnSchema("other", TSDataType.DOUBLE, ColumnCategory.FIELD), + ], + ) + with TsFileTableWriter(str(path), schema) as writer: + writer.write_dataframe( + pd.DataFrame( + { + "time": [0, 1, 0, 1, 0, 1], + "device": ["d0", "d0", "d1", "d1", "d2", "d2"], + "value": [0.0, 1.0, 10.0, 11.0, 20.0, 21.0], + "other": [100.0, np.nan, 110.0, 111.0, 120.0, 121.0], + } + ) + ) + return str(path) + + +def _effective_options(dataframe): + runtime = dataframe._runtime + return ( + runtime.prepared.max_entries, + runtime.catalog._descriptor_cache_size, + runtime.readers.max_open_files, + runtime.query_workers, + runtime.query_parallel_min_rows, + ) + + +def test_dataframe_options_override_environment(dataset_file, monkeypatch): + for option in _OPTIONS: + monkeypatch.setenv("TSFILE_DATAFRAME_" + option.upper(), "invalid") + options = dict(zip(_OPTIONS, (2, 3, 4, 1, 5))) + with TsFileDataFrame( + dataset_file, show_progress=False, use_index=True, **options + ) as dataframe: + assert _effective_options(dataframe) == (2, 3, 4, 1, 5) + with dataframe[[0, 1]] as subset: + assert subset._runtime is dataframe._runtime + assert _effective_options(subset) == _effective_options(dataframe) + + +def test_dataframe_snapshots_environment_for_each_runtime(dataset_file, monkeypatch): + for option in _OPTIONS: + monkeypatch.setenv("TSFILE_DATAFRAME_" + option.upper(), "2") + with TsFileDataFrame(dataset_file, show_progress=False, use_index=True) as first: + for option in _OPTIONS: + monkeypatch.setenv("TSFILE_DATAFRAME_" + option.upper(), "1") + with TsFileDataFrame( + dataset_file, show_progress=False, use_index=True + ) as second: + assert _effective_options(first) == (2, 2, 2, 2, 2) + assert _effective_options(second) == (1, 1, 1, 1, 1) + + +def test_dataframe_read_options_defaults(dataset_file, monkeypatch): + for option in _OPTIONS: + monkeypatch.delenv("TSFILE_DATAFRAME_" + option.upper(), raising=False) + monkeypatch.setattr("os.cpu_count", lambda: 2) + with TsFileDataFrame( + dataset_file, show_progress=False, use_index=True + ) as dataframe: + assert _effective_options(dataframe) == (4096, 4096, 16, 2, 8192) + + +@pytest.mark.parametrize("capacity", [0, 1, 2]) +@pytest.mark.parametrize("trust_index", [True, False]) +def test_small_caches_preserve_single_and_aligned_reads( + dataset_file, capacity, trust_index +): + with TsFileDataFrame( + dataset_file, + show_progress=False, + use_index=True, + trust_index=trust_index, + max_prepared_series=capacity, + descriptor_cache_size=capacity, + query_workers=2, + query_parallel_min_rows=1, + ) as dataframe: + names = [str(name) for name in dataframe.list_timeseries()] + baseline = dataframe.loc[0:1, names].values.copy() + expected = [] + for name in names: + device = int(name.split(".")[1][1:]) + values = np.array([device * 10.0, device * 10.0 + 1]) + if name.endswith(".other"): + values += 100 + if device == 0: + values[1] = np.nan + expected.append(values) + np.testing.assert_allclose(baseline, np.column_stack(expected), equal_nan=True) + for _ in range(2): + for column, name in enumerate(names): + with dataframe[name] as series: + np.testing.assert_allclose( + series[:], baseline[:, column], equal_nan=True + ) + assert dataframe._runtime.prepared.size <= capacity + assert len(dataframe._index._descriptor_cache) <= capacity + assert len(dataframe._index.series_shards._cache) <= capacity + np.testing.assert_allclose( + dataframe.loc[0:1, names].values, baseline, equal_nan=True + ) + assert dataframe._runtime.prepared.size <= capacity + + +@pytest.mark.parametrize("option", _OPTIONS) +@pytest.mark.parametrize("value", ["invalid", "-1", "1.5"]) +def test_invalid_environment_is_rejected_before_reading_files( + monkeypatch, option, value +): + env_name = "TSFILE_DATAFRAME_" + option.upper() + monkeypatch.setenv(env_name, value) + with pytest.raises(ValueError, match=env_name): + TsFileDataFrame("missing.tsfile", use_index=True) + + +@pytest.mark.parametrize("option", _OPTIONS[2:]) +def test_zero_requires_positive_read_limits(option): + with pytest.raises(ValueError, match=option): + TsFileDataFrame("missing.tsfile", use_index=True, **{option: 0}) + + +class _Prepared: + def __init__(self, locator_id, time_owner=None): + self.locator_id = locator_id + self.time_owner = time_owner + self.closed = False + + def close(self): + self.closed = True + + +class _Reader: + def __init__(self): + self.prepared = [] + + def prepare_series(self, locator, time_owner=None, trust_index=False): + assert time_owner is None or not time_owner.closed + prepared = _Prepared(locator, time_owner) + self.prepared.append(prepared) + return prepared + + +def _cache(monkeypatch, capacity): + cache = PreparedSeriesCache(SimpleNamespace(), max_entries=capacity) + monkeypatch.setattr(cache, "_locator_tuple", lambda _file, locator: locator) + return cache + + +def test_prepared_cache_evicts_lru_and_reprepares(monkeypatch): + cache = _cache(monkeypatch, 2) + reader = _Reader() + with cache.acquire(0, 0, reader) as first: + pass + with cache.acquire(0, 1, reader) as second: + pass + with cache.acquire(0, 0, reader) as reused: + assert reused is first + with cache.acquire(0, 2, reader): + assert second.closed + assert not first.closed + assert cache.size == 2 + with cache.acquire(0, 1, reader) as rebuilt: + assert rebuilt is not second + assert not rebuilt.closed + cache.close() + assert all(item.closed for item in reader.prepared) + + +def test_prepared_cache_preserves_access_order_when_leases_finish_out_of_order( + monkeypatch, +): + cache = _cache(monkeypatch, 2) + reader = _Reader() + with cache.acquire(0, 0, reader) as oldest: + with cache.acquire(0, 1, reader) as newer: + pass + with cache.acquire(0, 2, reader): + assert oldest.closed + assert not newer.closed + cache.close() + + +def test_prepared_cache_restores_unpinned_owner_in_access_order(monkeypatch): + cache = _cache(monkeypatch, 3) + reader = _Reader() + with cache.acquire(0, 0, reader) as owner: + with cache.acquire(0, 1, reader, time_owner=owner) as dependent: + pass + with cache.acquire(0, 2, reader) as newer: + pass + with cache.acquire(0, 3, reader): + assert dependent.closed + assert not owner.closed + with cache.acquire(0, 4, reader): + assert owner.closed + assert not newer.closed + cache.close() + + +def test_prepared_cache_bounds_idle_queue_during_repeated_cache_hits(monkeypatch): + cache = _cache(monkeypatch, 2) + reader = _Reader() + for _ in range(256): + with cache.acquire(0, 0, reader) as prepared: + assert not prepared.closed + # Lazy invalidation must not retain one heap record per cache hit. + assert len(cache._idle_heap) <= 128 + assert len(reader.prepared) == 1 + cache.close() + assert prepared.closed + + +@pytest.mark.parametrize("capacity", [0, 1, 32]) +@pytest.mark.parametrize("aligned", [False, True]) +def test_wide_prepared_query_does_not_repeatedly_scan_pinned_entries( + monkeypatch, capacity, aligned +): + class CountedEntries(OrderedDict): + visited = 0 + + def items(self): + for item in super().items(): + self.visited += 1 + yield item + + def values(self): + for entry in super().values(): + self.visited += 1 + yield entry + + cache = _cache(monkeypatch, capacity) + entries = CountedEntries() + monkeypatch.setattr(cache, "_entries", entries) + reader = _Reader() + width = 128 + with ExitStack() as stack: + owner = None + for locator in range(width): + prepared = stack.enter_context( + cache.acquire(0, locator, reader, time_owner=owner) + ) + if aligned and owner is None: + owner = prepared + assert cache.size == width + assert all(not prepared.closed for prepared in reader.prepared) + assert cache.size <= capacity + cache.close() + assert all(prepared.closed for prepared in reader.prepared) + assert entries.visited <= 8 * width + + +def test_prepared_cache_zero_keeps_active_entries_and_shared_owner(monkeypatch): + cache = _cache(monkeypatch, 0) + reader = _Reader() + with cache.acquire(0, 0, reader) as owner: + with cache.acquire(0, 1, reader, time_owner=owner) as value: + assert value.time_owner is owner + assert not owner.closed + assert not value.closed + assert cache.size == 2 + assert value.closed + assert not owner.closed + assert owner.closed + assert cache.size == 0 + cache.close() + + +def test_prepared_cache_keeps_owner_until_dependent_is_evicted(monkeypatch): + cache = _cache(monkeypatch, 2) + reader = _Reader() + with cache.acquire(0, 0, reader) as owner: + with cache.acquire(0, 1, reader, time_owner=owner) as value: + pass + with cache.acquire(0, 2, reader): + assert value.closed + assert not owner.closed + with cache.acquire(0, 3, reader): + assert owner.closed + cache.close() + + +def test_prepared_cache_single_flight_and_active_queries_survive_pressure(monkeypatch): + cache = _cache(monkeypatch, 1) + reader = _Reader() + preparing = threading.Event() + finish_preparing = threading.Event() + acquired = threading.Barrier(3, timeout=5) + release = threading.Event() + original = reader.prepare_series + + def prepare(locator, time_owner=None, trust_index=False): + preparing.set() + assert finish_preparing.wait(timeout=5) + return original(locator, time_owner, trust_index=trust_index) + + monkeypatch.setattr(reader, "prepare_series", prepare) + + def query(): + with cache.acquire(0, 0, reader) as prepared: + acquired.wait() + assert release.wait(timeout=5) + assert not prepared.closed + return prepared + + with ThreadPoolExecutor(max_workers=2) as executor: + first = executor.submit(query) + assert preparing.wait(timeout=5) + second = executor.submit(query) + finish_preparing.set() + try: + acquired.wait() + assert len(reader.prepared) == 1 + with cache.acquire(0, 1, reader) as other: + assert not reader.prepared[0].closed + assert not other.closed + assert other.closed + finally: + release.set() + assert first.result(timeout=5) is second.result(timeout=5) + cache.close() + + +def test_prepared_cache_close_waits_for_active_lease(monkeypatch): + cache = _cache(monkeypatch, 0) + reader = _Reader() + closing = threading.Event() + closed = threading.Event() + + def close(): + closing.set() + cache.close() + closed.set() + + with ThreadPoolExecutor(max_workers=1) as executor: + with cache.acquire(0, 0, reader) as prepared: + future = executor.submit(close) + assert closing.wait(timeout=5) + assert not closed.wait(timeout=0.05) + assert not prepared.closed + future.result(timeout=5) + assert prepared.closed + assert closed.is_set() + + +def test_failed_preparation_can_be_retried(monkeypatch): + cache = _cache(monkeypatch, 1) + reader = _Reader() + original = reader.prepare_series + + def fail(*_args, **_kwargs): + raise RuntimeError("prepare failed") + + monkeypatch.setattr(reader, "prepare_series", fail) + with pytest.raises(RuntimeError, match="prepare failed"): + with cache.acquire(0, 0, reader): + pass + monkeypatch.setattr(reader, "prepare_series", original) + with cache.acquire(0, 0, reader) as prepared: + assert not prepared.closed + cache.close() + + +def test_cache_close_waits_for_eviction_cleanup(monkeypatch): + cache = _cache(monkeypatch, 0) + reader = _Reader() + freeing = threading.Event() + finish_freeing = threading.Event() + closing = threading.Event() + closed = threading.Event() + original_prepare = reader.prepare_series + + def prepare(*args, **kwargs): + prepared = original_prepare(*args, **kwargs) + original_close = prepared.close + + def free(): + freeing.set() + assert finish_freeing.wait(timeout=5) + original_close() + + prepared.close = free + return prepared + + monkeypatch.setattr(reader, "prepare_series", prepare) + + def query(): + with cache.acquire(0, 0, reader): + pass + + def close(): + closing.set() + cache.close() + closed.set() + + with ThreadPoolExecutor(max_workers=2) as executor: + query_future = executor.submit(query) + assert freeing.wait(timeout=5) + close_future = executor.submit(close) + try: + assert closing.wait(timeout=5) + assert not closed.wait(timeout=0.05) + finally: + finish_freeing.set() + query_future.result(timeout=5) + close_future.result(timeout=5) + assert all(item.closed for item in reader.prepared) + + +def test_cache_close_waits_for_preparation_and_releases_rejected_result(monkeypatch): + cache = _cache(monkeypatch, 1) + reader = _Reader() + preparing = threading.Event() + finish_preparing = threading.Event() + closing = threading.Event() + closed = threading.Event() + original = reader.prepare_series + + def prepare(*args, **kwargs): + preparing.set() + assert finish_preparing.wait(timeout=5) + return original(*args, **kwargs) + + monkeypatch.setattr(reader, "prepare_series", prepare) + + def query(): + with cache.acquire(0, 0, reader): + pytest.fail("closing cache must reject a newly prepared result") + + def close(): + closing.set() + cache.close() + closed.set() + + with ThreadPoolExecutor(max_workers=2) as executor: + query_future = executor.submit(query) + assert preparing.wait(timeout=5) + close_future = executor.submit(close) + try: + assert closing.wait(timeout=5) + assert not closed.wait(timeout=0.05) + finally: + finish_preparing.set() + with pytest.raises(RuntimeError, match="closed"): + query_future.result(timeout=5) + close_future.result(timeout=5) + assert all(item.closed for item in reader.prepared) + assert cache.size == 0 + + +@pytest.mark.parametrize( + "option", + [ + "max_prepared_series", + "descriptor_cache_size", + "max_open_files", + "query_workers", + "query_parallel_min_rows", + ], +) +@pytest.mark.parametrize("value", [-1, 1.5, True, "32"]) +def test_dataframe_rejects_invalid_options_before_reading_files(option, value): + with pytest.raises((ValueError, TypeError), match=option): + TsFileDataFrame("missing.tsfile", use_index=True, **{option: value}) + + +def test_dataframe_rejects_ignored_options_without_index(): + with pytest.raises(ValueError, match="use_index=True"): + TsFileDataFrame("missing.tsfile", max_prepared_series=32) diff --git a/python/tsfile/dataset/_config.py b/python/tsfile/dataset/_config.py new file mode 100644 index 000000000..0fbdd22ea --- /dev/null +++ b/python/tsfile/dataset/_config.py @@ -0,0 +1,66 @@ +# 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. + +"""Resolve per-runtime read options before opening dataset resources.""" + +import operator +import os + +DEFAULT_CACHE_SIZE = 4096 + + +def resolve_read_options( + *, + max_prepared_series=None, + descriptor_cache_size=None, + max_open_files=None, + query_workers=None, + query_parallel_min_rows=None, +): + """Use explicit values, then environment variables, then built-in defaults.""" + defaults = { + "max_prepared_series": (max_prepared_series, DEFAULT_CACHE_SIZE, 0), + "descriptor_cache_size": (descriptor_cache_size, DEFAULT_CACHE_SIZE, 0), + "max_open_files": (max_open_files, 16, 1), + "query_workers": (query_workers, min(4, os.cpu_count() or 1), 1), + "query_parallel_min_rows": (query_parallel_min_rows, 8192, 1), + } + options = {} + for name, (value, default, minimum) in defaults.items(): + source = name + if value is None: + source = "TSFILE_DATAFRAME_" + name.upper() + raw = os.environ.get(source) + if raw is None: + value = default + else: + try: + value = int(raw) + except ValueError as exc: + raise ValueError(f"{source} must be an integer") from exc + else: + if isinstance(value, bool): + raise TypeError(f"{source} must be an integer, not bool") + try: + value = operator.index(value) + except TypeError as exc: + raise TypeError(f"{source} must be an integer") from exc + if value < minimum: + requirement = "non-negative" if minimum == 0 else "positive" + raise ValueError(f"{source} must be {requirement}") + options[name] = value + return options diff --git a/python/tsfile/dataset/dataframe.py b/python/tsfile/dataset/dataframe.py index 7ccdc68f9..893419d96 100644 --- a/python/tsfile/dataset/dataframe.py +++ b/python/tsfile/dataset/dataframe.py @@ -30,6 +30,7 @@ import numpy as np from .formatting import format_dataframe_table +from ._config import resolve_read_options from .metadata import ( MODEL_TABLE, MODEL_TREE, @@ -772,19 +773,51 @@ def __getitem__(self, key) -> AlignedTimeseries: class TsFileDataFrame: - """Lazy-loaded unified numeric dataset view over multiple TsFile shards.""" + """Lazy-loaded unified numeric dataset view over multiple TsFile shards. + + With ``use_index=True``, ``trust_index=True`` (the default) assumes the + indexed TsFiles are immutable and skips source generation checks. Set + ``trust_index=False`` to check file sizes and modification times when + loading the index and acquiring readers. Index format and locator bounds + are always checked. ``trust_index`` has no effect when ``use_index=False``. + + Keyword-only resource limits apply to ``use_index=True``. Each limit uses + its explicit value, or the matching ``TSFILE_DATAFRAME_