fix(cpp): apply residual row offsets in non-aligned page decoding (#1007)
diff --git a/cpp/src/reader/chunk_reader.cc b/cpp/src/reader/chunk_reader.cc
index abeff83..d713570 100644
--- a/cpp/src/reader/chunk_reader.cc
+++ b/cpp/src/reader/chunk_reader.cc
@@ -315,7 +315,7 @@
}
int ChunkReader::decode_cur_page_data(TsBlock*& ret_tsblock, Filter* filter,
- PageArena& pa) {
+ PageArena& pa, int* row_offset) {
int ret = E_OK;
// Step 1: make sure we load the whole page data in @in_stream_
@@ -396,8 +396,8 @@
// ret = decode_tv_buf_into_tsblock(time_buf, value_buf, time_buf_size,
// value_buf_size, ret_tsblock,
// filter);
- ret = decode_tv_buf_into_tsblock_by_datatype(time_in_, value_in_,
- ret_tsblock, filter, &pa);
+ ret = decode_tv_buf_into_tsblock_by_datatype(
+ time_in_, value_in_, ret_tsblock, filter, &pa, row_offset);
// if we return during @decode_tv_buf_into_tsblock, we should keep
// @uncompressed_buf_ valid until all TV pairs are decoded.
if (ret != E_OVERFLOW) {
@@ -414,29 +414,32 @@
return ret;
}
-#define DECODE_TYPED_TV_INTO_TSBLOCK(CppType, ReadType, time_in, value_in, \
- row_appender) \
- do { \
- int64_t time = 0; \
- CppType value; \
- while (time_decoder_->has_remaining(time_in)) { \
- ASSERT(value_decoder_->has_remaining(value_in)); \
- if (UNLIKELY(!row_appender.add_row())) { \
- ret = E_OVERFLOW; \
- break; \
- } else if (RET_FAIL(time_decoder_->read_int64(time, time_in))) { \
- } else if (RET_FAIL(value_decoder_->read_##ReadType(value, \
- value_in))) { \
- } else if (filter != nullptr && !filter->satisfy(time, value)) { \
- row_appender.backoff_add_row(); \
- continue; \
- } else { \
- /*std::cout << "decoder: time=" << time << ", value=" << value \
- * << std::endl;*/ \
- row_appender.append(0, (char*)&time, sizeof(time)); \
- row_appender.append(1, (char*)&value, sizeof(value)); \
- } \
- } \
+#define DECODE_TYPED_TV_INTO_TSBLOCK(CppType, ReadType, time_in, value_in, \
+ row_appender) \
+ do { \
+ int64_t time = 0; \
+ CppType value; \
+ while (time_decoder_->has_remaining(time_in)) { \
+ ASSERT(value_decoder_->has_remaining(value_in)); \
+ if (UNLIKELY(!row_appender.add_row())) { \
+ ret = E_OVERFLOW; \
+ break; \
+ } else if (RET_FAIL(time_decoder_->read_int64(time, time_in))) { \
+ } else if (RET_FAIL(value_decoder_->read_##ReadType(value, \
+ value_in))) { \
+ } else if (filter != nullptr && !filter->satisfy(time, value)) { \
+ row_appender.backoff_add_row(); \
+ continue; \
+ } else { \
+ if (row_offset > 0) { \
+ --row_offset; \
+ row_appender.backoff_add_row(); \
+ continue; \
+ } \
+ row_appender.append(0, (char*)&time, sizeof(time)); \
+ row_appender.append(1, (char*)&value, sizeof(value)); \
+ } \
+ } \
} while (false)
int ChunkReader::i32_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in,
@@ -467,19 +470,19 @@
}
int ChunkReader::i32_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in,
- RowAppender& row_appender,
- Filter* filter) {
+ RowAppender& row_appender, Filter* filter,
+ int& row_offset) {
int ret = E_OK;
const int BATCH = 129;
int64_t times[BATCH];
int32_t values[BATCH];
while (time_decoder_->has_remaining(time_in)) {
- // Cap each pass to what the appender can still hold; the old
- // "remaining < BATCH → OVERFLOW" check made progress impossible on
- // TsBlocks with capacity below BATCH.
- int eff_batch =
- std::min(BATCH, static_cast<int>(row_appender.remaining()));
+ // Skipped rows consume no output capacity. Bound the batch so every
+ // accepted row fits after the remaining offset has been consumed.
+ int eff_batch = std::min(
+ BATCH,
+ std::max(static_cast<int>(row_appender.remaining()), row_offset));
if (eff_batch <= 0) {
ret = E_OVERFLOW;
break;
@@ -543,6 +546,10 @@
!filter->satisfy(times[i], (int64_t)values[i])) {
continue;
}
+ if (row_offset > 0) {
+ --row_offset;
+ continue;
+ }
if (UNLIKELY(!row_appender.add_row())) {
ret = E_OVERFLOW;
break;
@@ -556,16 +563,17 @@
}
int ChunkReader::i64_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in,
- RowAppender& row_appender,
- Filter* filter) {
+ RowAppender& row_appender, Filter* filter,
+ int& row_offset) {
int ret = E_OK;
const int BATCH = 129;
int64_t times[BATCH];
int64_t values[BATCH];
while (time_decoder_->has_remaining(time_in)) {
- int eff_batch =
- std::min(BATCH, static_cast<int>(row_appender.remaining()));
+ int eff_batch = std::min(
+ BATCH,
+ std::max(static_cast<int>(row_appender.remaining()), row_offset));
if (eff_batch <= 0) {
ret = E_OVERFLOW;
break;
@@ -629,6 +637,10 @@
!filter->satisfy(times[i], values[i])) {
continue;
}
+ if (row_offset > 0) {
+ --row_offset;
+ continue;
+ }
if (UNLIKELY(!row_appender.add_row())) {
ret = E_OVERFLOW;
break;
@@ -644,15 +656,16 @@
int ChunkReader::float_DECODE_TV_BATCH(ByteStream& time_in,
ByteStream& value_in,
RowAppender& row_appender,
- Filter* filter) {
+ Filter* filter, int& row_offset) {
int ret = E_OK;
const int BATCH = 129;
int64_t times[BATCH];
float values[BATCH];
while (time_decoder_->has_remaining(time_in)) {
- int eff_batch =
- std::min(BATCH, static_cast<int>(row_appender.remaining()));
+ int eff_batch = std::min(
+ BATCH,
+ std::max(static_cast<int>(row_appender.remaining()), row_offset));
if (eff_batch <= 0) {
ret = E_OVERFLOW;
break;
@@ -712,6 +725,10 @@
if (filter != nullptr && !block_all_pass && !time_mask[i]) {
continue;
}
+ if (row_offset > 0) {
+ --row_offset;
+ continue;
+ }
if (UNLIKELY(!row_appender.add_row())) {
ret = E_OVERFLOW;
break;
@@ -727,15 +744,16 @@
int ChunkReader::double_DECODE_TV_BATCH(ByteStream& time_in,
ByteStream& value_in,
RowAppender& row_appender,
- Filter* filter) {
+ Filter* filter, int& row_offset) {
int ret = E_OK;
const int BATCH = 129;
int64_t times[BATCH];
double values[BATCH];
while (time_decoder_->has_remaining(time_in)) {
- int eff_batch =
- std::min(BATCH, static_cast<int>(row_appender.remaining()));
+ int eff_batch = std::min(
+ BATCH,
+ std::max(static_cast<int>(row_appender.remaining()), row_offset));
if (eff_batch <= 0) {
ret = E_OVERFLOW;
break;
@@ -795,6 +813,10 @@
if (filter != nullptr && !block_all_pass && !time_mask[i]) {
continue;
}
+ if (row_offset > 0) {
+ --row_offset;
+ continue;
+ }
if (UNLIKELY(!row_appender.add_row())) {
ret = E_OVERFLOW;
break;
@@ -807,11 +829,9 @@
return ret;
}
-int ChunkReader::STRING_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in,
- ByteStream& value_in,
- RowAppender& row_appender,
- PageArena& pa,
- Filter* filter) {
+int ChunkReader::STRING_DECODE_TYPED_TV_INTO_TSBLOCK(
+ ByteStream& time_in, ByteStream& value_in, RowAppender& row_appender,
+ PageArena& pa, Filter* filter, int& row_offset) {
int ret = E_OK;
int64_t time = 0;
common::String value;
@@ -825,6 +845,10 @@
} else if (filter != nullptr && !filter->satisfy(time, value)) {
row_appender.backoff_add_row();
continue;
+ } else if (row_offset > 0) {
+ --row_offset;
+ row_appender.backoff_add_row();
+ continue;
} else {
row_appender.append(0, (char*)&time, sizeof(time));
row_appender.append(1, value.buf_, value.len_);
@@ -833,12 +857,13 @@
return ret;
}
-int ChunkReader::decode_tv_buf_into_tsblock_by_datatype(ByteStream& time_in,
- ByteStream& value_in,
- TsBlock* ret_tsblock,
- Filter* filter,
- common::PageArena* pa) {
+int ChunkReader::decode_tv_buf_into_tsblock_by_datatype(
+ ByteStream& time_in, ByteStream& value_in, TsBlock* ret_tsblock,
+ Filter* filter, common::PageArena* pa, int* remaining_offset) {
int ret = E_OK;
+ int unused_offset = 0;
+ int& row_offset =
+ remaining_offset == nullptr ? unused_offset : *remaining_offset;
RowAppender row_appender(ret_tsblock);
switch (chunk_header_.data_type_) {
case common::BOOLEAN:
@@ -847,33 +872,36 @@
break;
case common::DATE:
case common::INT32:
- ret =
- i32_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter);
+ ret = i32_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter,
+ row_offset);
break;
case TIMESTAMP:
case common::INT64:
- ret =
- i64_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter);
+ ret = i64_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter,
+ row_offset);
break;
case common::FLOAT:
ret = float_DECODE_TV_BATCH(time_in_, value_in_, row_appender,
- filter);
+ filter, row_offset);
break;
case common::DOUBLE:
ret = double_DECODE_TV_BATCH(time_in_, value_in_, row_appender,
- filter);
+ filter, row_offset);
break;
case common::TEXT:
case common::BLOB:
case common::STRING:
ret = STRING_DECODE_TYPED_TV_INTO_TSBLOCK(
- time_in, value_in, row_appender, *pa, filter);
+ time_in, value_in, row_appender, *pa, filter, row_offset);
break;
default:
ret = E_NOT_SUPPORT;
ASSERT(false);
}
- if (ret_tsblock->get_row_count() == 0 && ret == E_OK) {
+ // The offset-aware iterator must keep scanning after a page whose
+ // matching rows were all consumed by the offset (or time filter).
+ if (remaining_offset == nullptr && ret_tsblock->get_row_count() == 0 &&
+ ret == E_OK) {
ret = E_NO_MORE_DATA;
}
return ret;
@@ -917,8 +945,8 @@
}
if (prev_page_not_finish()) {
- ret = decode_tv_buf_into_tsblock_by_datatype(time_in_, value_in_,
- ret_tsblock, filter, &pa);
+ ret = decode_tv_buf_into_tsblock_by_datatype(
+ time_in_, value_in_, ret_tsblock, filter, &pa, &row_offset);
if (ret == E_OVERFLOW) {
ret = E_OK;
} else {
@@ -953,7 +981,7 @@
}
if (IS_SUCC(ret)) {
- ret = decode_cur_page_data(ret_tsblock, filter, pa);
+ ret = decode_cur_page_data(ret_tsblock, filter, pa, &row_offset);
}
return ret;
}
diff --git a/cpp/src/reader/chunk_reader.h b/cpp/src/reader/chunk_reader.h
index b543011..e49f727 100644
--- a/cpp/src/reader/chunk_reader.h
+++ b/cpp/src/reader/chunk_reader.h
@@ -91,7 +91,7 @@
bool cur_page_fully_satisfies_filter(Filter* filter);
int skip_cur_page();
int decode_cur_page_data(common::TsBlock*& ret_tsblock, Filter* filter,
- common::PageArena& pa);
+ common::PageArena& pa, int* row_offset = nullptr);
bool prev_page_not_finish() const {
return (time_decoder_ && time_decoder_->has_remaining(time_in_)) ||
time_in_.has_remaining();
@@ -101,30 +101,33 @@
common::ByteStream& value_in,
common::TsBlock* ret_tsblock,
Filter* filter,
- common::PageArena* pa = nullptr);
+ common::PageArena* pa = nullptr,
+ int* remaining_offset = nullptr);
int i32_DECODE_TYPED_TV_INTO_TSBLOCK(common::ByteStream& time_in,
common::ByteStream& value_in,
common::RowAppender& row_appender,
Filter* filter);
int i32_DECODE_TV_BATCH(common::ByteStream& time_in,
common::ByteStream& value_in,
- common::RowAppender& row_appender, Filter* filter);
+ common::RowAppender& row_appender, Filter* filter,
+ int& row_offset);
int i64_DECODE_TV_BATCH(common::ByteStream& time_in,
common::ByteStream& value_in,
- common::RowAppender& row_appender, Filter* filter);
+ common::RowAppender& row_appender, Filter* filter,
+ int& row_offset);
int float_DECODE_TV_BATCH(common::ByteStream& time_in,
common::ByteStream& value_in,
- common::RowAppender& row_appender,
- Filter* filter);
+ common::RowAppender& row_appender, Filter* filter,
+ int& row_offset);
int double_DECODE_TV_BATCH(common::ByteStream& time_in,
common::ByteStream& value_in,
common::RowAppender& row_appender,
- Filter* filter);
+ Filter* filter, int& row_offset);
int STRING_DECODE_TYPED_TV_INTO_TSBLOCK(common::ByteStream& time_in,
common::ByteStream& value_in,
common::RowAppender& row_appender,
common::PageArena& pa,
- Filter* filter);
+ Filter* filter, int& row_offset);
private:
RandomAccessReadFile* read_file_;
diff --git a/cpp/test/reader/prepared_series_test.cc b/cpp/test/reader/prepared_series_test.cc
index cd1389f..d90101f 100644
--- a/cpp/test/reader/prepared_series_test.cc
+++ b/cpp/test/reader/prepared_series_test.cc
@@ -22,6 +22,7 @@
#include <gtest/gtest.h>
#include <sys/stat.h>
+#include <algorithm>
#include <cstdio>
#include <vector>
@@ -32,22 +33,29 @@
#include "reader/table_result_set.h"
#include "reader/tsfile_reader.h"
#include "writer/tsfile_table_writer.h"
+#include "writer/tsfile_writer.h"
namespace storage {
namespace {
class PagePointGuard {
public:
- explicit PagePointGuard(uint32_t page_points)
- : saved_(common::g_config_value_.page_writer_max_point_num_) {
+ explicit PagePointGuard(uint32_t page_points, uint32_t page_bytes = 0)
+ : saved_(common::g_config_value_.page_writer_max_point_num_),
+ saved_bytes_(common::g_config_value_.page_writer_max_memory_bytes_) {
common::g_config_value_.page_writer_max_point_num_ = page_points;
+ if (page_bytes != 0) {
+ common::g_config_value_.page_writer_max_memory_bytes_ = page_bytes;
+ }
}
~PagePointGuard() {
common::g_config_value_.page_writer_max_point_num_ = saved_;
+ common::g_config_value_.page_writer_max_memory_bytes_ = saved_bytes_;
}
private:
uint32_t saved_;
+ uint32_t saved_bytes_;
};
class PreparedSeriesBatchTest : public ::testing::Test {
@@ -59,6 +67,7 @@
file_name_ = std::string("prepared_series_batch_test_") +
(test_info == nullptr ? "unknown" : test_info->name()) +
".tsfile";
+ std::replace(file_name_.begin(), file_name_.end(), '/', '_');
std::remove(file_name_.c_str());
}
@@ -111,6 +120,173 @@
std::string file_name_ = "prepared_series_batch_test.tsfile";
};
+class NonAlignedPreparedSeriesOffsetTest
+ : public PreparedSeriesBatchTest,
+ public ::testing::WithParamInterface<uint32_t> {};
+
+TEST_P(NonAlignedPreparedSeriesOffsetTest, AppliesRowOffset) {
+ // Include pages larger than the 65536-row output block so the decoder
+ // resumes an unfinished page on the next read.
+ PagePointGuard guard(GetParam(), 16 * 1024 * 1024);
+ const std::string device = "root.offset";
+ const std::vector<std::string> names = {"boolean", "int32", "int64",
+ "float", "double", "string"};
+ const std::vector<common::TSDataType> types = {
+ common::BOOLEAN, common::INT32, common::INT64,
+ common::FLOAT, common::DOUBLE, common::STRING};
+ auto schemas = std::make_shared<std::vector<MeasurementSchema>>();
+ TsFileWriter writer;
+ ASSERT_EQ(common::E_OK, writer.open(file_name_));
+ for (size_t i = 0; i < names.size(); ++i) {
+ schemas->emplace_back(names[i], types[i], common::PLAIN,
+ common::UNCOMPRESSED);
+ ASSERT_EQ(common::E_OK,
+ writer.register_timeseries(device, schemas->back()));
+ }
+ // Flush halfway through to exercise both page and chunk skipping.
+ for (int start : {0, 70000}) {
+ Tablet tablet(device, schemas, 70000);
+ for (int row = 0; row < 70000; ++row) {
+ const int value = start + row;
+ ASSERT_EQ(common::E_OK, tablet.add_timestamp(row, value));
+ ASSERT_EQ(common::E_OK, tablet.add_value(row, 0u, value % 2 == 0));
+ ASSERT_EQ(common::E_OK, tablet.add_value(row, 1u, int32_t(value)));
+ ASSERT_EQ(common::E_OK, tablet.add_value(row, 2u, int64_t(value)));
+ ASSERT_EQ(common::E_OK,
+ tablet.add_value(row, 3u, float(value) + 0.5f));
+ ASSERT_EQ(common::E_OK,
+ tablet.add_value(row, 4u, double(value) + 0.5));
+ const std::string text = std::to_string(value);
+ ASSERT_EQ(common::E_OK, tablet.add_value(row, 5u, text.c_str()));
+ }
+ ASSERT_EQ(common::E_OK, writer.write_tablet(tablet));
+ ASSERT_EQ(common::E_OK, writer.flush());
+ }
+ ASSERT_EQ(common::E_OK, writer.close());
+
+ TsFileReader reader;
+ ASSERT_EQ(common::E_OK, reader.open(file_name_));
+ FileGeneration generation;
+ generation.mapped_index_identity = 1;
+ generation.file_id = 0;
+ struct stat file_stat {};
+ ASSERT_EQ(0, stat(file_name_.c_str(), &file_stat));
+ generation.file_size = static_cast<uint64_t>(file_stat.st_size);
+ generation.file_fingerprint = 0;
+
+ struct Window {
+ int64_t start;
+ int64_t end;
+ int offset;
+ int limit;
+ };
+ const std::vector<Window> windows = {
+ {0, 139999, 30, 10}, // Offset inside the first page.
+ {0, 139999, 10003, 1}, // Output capacity below the decode batch size.
+ {0, 139999, 70003, 7}, // Whole chunk plus a residual.
+ {0, 139999, 13, 65537}, // Resumes decoding a large page.
+ {0, 139999, 69995, 10}, // Window crossing a chunk boundary.
+ {0, 139999, 139990, -1}, // Unlimited tail.
+ {0, 139999, 140001, 5}, // Offset beyond the data.
+ {99, 150, 5, 7}, // Count only rows passing the time filter.
+ {99, 150, 60, 7}, // Offset beyond the filtered rows.
+ {9998, 10005, 3, 4}, // Filter and offset crossing a page boundary.
+ {0, 139999, 30, 0},
+ };
+ auto metadata = reader.get_timeseries_metadata();
+ ASSERT_EQ(1U, metadata.size());
+ ASSERT_EQ(names.size(), metadata.begin()->second.size());
+ for (const auto& device_entry : metadata) {
+ for (const auto& index : device_entry.second) {
+ auto* series_index = dynamic_cast<TimeseriesIndex*>(index.get());
+ ASSERT_NE(nullptr, series_index);
+ PreparedLocator locator;
+ locator.layout = 0;
+ locator.value_metadata_offset = series_index->get_metadata_offset();
+ locator.value_metadata_length = series_index->get_metadata_length();
+ std::shared_ptr<PreparedSeries> prepared;
+ ASSERT_EQ(common::E_OK,
+ reader.prepare_series(generation, locator, prepared));
+ for (const auto& window : windows) {
+ SCOPED_TRACE(index->get_measurement_name().to_std_string());
+ SCOPED_TRACE(window.offset);
+ ResultSet* result = nullptr;
+ ASSERT_EQ(
+ common::E_OK,
+ reader.query_prepared(prepared, window.start, window.end,
+ window.offset, window.limit, result));
+ auto* table = dynamic_cast<TableResultSet*>(result);
+ ASSERT_NE(nullptr, table);
+ const int available = std::max(
+ 0, int(window.end - window.start + 1) - window.offset);
+ const int expected_count =
+ window.limit < 0 ? available
+ : std::min(available, window.limit);
+ int count = 0;
+ common::TsBlock* block = nullptr;
+ int ret = common::E_OK;
+ while ((ret = table->get_next_tsblock(block)) == common::E_OK) {
+ common::RowIterator rows(block);
+ while (rows.has_next()) {
+ const int64_t expected =
+ window.start + window.offset + count;
+ uint32_t len = 0;
+ bool is_null = false;
+ const char* time = rows.read(0, &len, &is_null);
+ ASSERT_FALSE(is_null);
+ ASSERT_EQ(expected,
+ *reinterpret_cast<const int64_t*>(time));
+ const char* value = rows.read(1, &len, &is_null);
+ ASSERT_FALSE(is_null);
+ switch (index->get_data_type()) {
+ case common::BOOLEAN:
+ EXPECT_EQ(
+ expected % 2 == 0,
+ *reinterpret_cast<const bool*>(value));
+ break;
+ case common::INT32:
+ EXPECT_EQ(
+ expected,
+ *reinterpret_cast<const int32_t*>(value));
+ break;
+ case common::INT64:
+ EXPECT_EQ(
+ expected,
+ *reinterpret_cast<const int64_t*>(value));
+ break;
+ case common::FLOAT:
+ EXPECT_FLOAT_EQ(
+ expected + 0.5f,
+ *reinterpret_cast<const float*>(value));
+ break;
+ case common::DOUBLE:
+ EXPECT_DOUBLE_EQ(
+ expected + 0.5,
+ *reinterpret_cast<const double*>(value));
+ break;
+ case common::STRING:
+ EXPECT_EQ(std::to_string(expected),
+ std::string(value, len));
+ break;
+ default:
+ FAIL() << "Unexpected type";
+ }
+ ++count;
+ rows.next();
+ }
+ }
+ EXPECT_EQ(common::E_NO_MORE_DATA, ret);
+ EXPECT_EQ(expected_count, count);
+ reader.destroy_query_data_set(result);
+ }
+ }
+ }
+ EXPECT_EQ(common::E_OK, reader.close());
+}
+
+INSTANTIATE_TEST_SUITE_P(PageSizes, NonAlignedPreparedSeriesOffsetTest,
+ ::testing::Values(10000U, 100000U));
+
TEST_F(PreparedSeriesBatchTest,
PreparedQueryReturnsDirectTableResultSetBatches) {
write_nullable_table();
diff --git a/python/tests/test_tsfile_dataset.py b/python/tests/test_tsfile_dataset.py
index 02baa9b..964dde6 100644
--- a/python/tests/test_tsfile_dataset.py
+++ b/python/tests/test_tsfile_dataset.py
@@ -2329,6 +2329,38 @@
writer.close()
+@pytest.mark.parametrize("dtype", [TSDataType.INT32, TSDataType.DOUBLE])
+def test_dataset_tree_model_row_slices_apply_page_offset(
+ tmp_path, dataframe_use_index, dtype
+):
+ paths = [tmp_path / "part0.tsfile", tmp_path / "part1.tsfile"]
+ for start, path in zip((0, 40), paths):
+ _write_tree_rows(
+ path, {"root.offset": [("value", dtype)]}, t_start=start, t_count=40
+ )
+ expected = np.arange(80, dtype=np.float64)
+ if dtype == TSDataType.DOUBLE:
+ expected += 0.5
+
+ # Check both the initial build and reopening the persisted index.
+ for _ in range(2):
+ with TsFileDataFrame(
+ [str(path) for path in paths],
+ show_progress=False,
+ use_index=dataframe_use_index,
+ ) as dataframe:
+ series = dataframe["root.offset.value"]
+ assert series[30] == expected[30]
+ for window in (
+ slice(30, 40),
+ slice(35, 55),
+ slice(43, 47),
+ slice(-10, None),
+ slice(30, 40, 3),
+ ):
+ np.testing.assert_array_equal(series[window], expected[window])
+
+
def test_dataset_tree_model_merges_identical_structure_across_files(
tmp_path, dataframe_use_index
):