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: 5 additions & 2 deletions cpp/src/cwrapper/tsfile_cwrapper.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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;
}
Expand All @@ -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;
Expand Down
10 changes: 7 additions & 3 deletions cpp/src/cwrapper/tsfile_cwrapper.h
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*/
Expand Down
11 changes: 7 additions & 4 deletions cpp/src/file/tsfile_io_reader.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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<PreparedSeries> candidate =
std::make_shared<PreparedSeries>(generation, locator);
TimeseriesIndex* value_index = nullptr;
Expand Down
5 changes: 4 additions & 1 deletion cpp/src/reader/prepared_series.h
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
61 changes: 61 additions & 0 deletions python/README-zh.md
Original file line number Diff line number Diff line change
Expand Up @@ -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。
Expand Down
67 changes: 67 additions & 0 deletions python/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.
Expand Down
98 changes: 86 additions & 12 deletions python/tests/test_dataset_index.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand All @@ -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:
Expand All @@ -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):
Expand Down
Loading
Loading