From 4a11bc81d4b4bbad5ce1a0991354374b47cc15db Mon Sep 17 00:00:00 2001 From: ColinLee Date: Tue, 22 Sep 2026 17:55:41 +0800 Subject: [PATCH 1/4] fix(cpp): keep aligned value columns row-aligned when a row omits measurements Every value column of an aligned chunk group has to consume exactly one row per time column row. write_record_aligned() / write_tablet_aligned() only advanced the columns that the incoming record/tablet carried, so a measurement missing from a row left its column one row short: its not-null bitmap started at the wrong row index and the values were paired with the earliest timestamps of the page on read, and the statistics recomputed after recovery described the shifted rows. - create the value chunk writer when the measurement is registered on an aligned device, so every registered measurement takes part from row 0 - write NULL rows for the measurements a record/tablet does not carry (Java: AlignedChunkGroupWriterImpl#write -> writeEmptyDataInOneRow) - reject registering a new measurement once rows have been written, which would need backfilled rows/pages; Java does not allow expanding an aligned device either - ValuePageWriter::write_null_rows() / ValueChunkWriter::write_null_batch() advance the column by NULL rows and keep page boundaries in step with the time column - a record repeating a measurement now advances that column once (last point wins) instead of running ahead of the time column Tests: TsFileWriterTest.AlignedRecordMissingMeasurementsStayRowAligned, AlignedRecordMissingMeasurementsAcrossPages, AlignedTabletMissingColumnStaysRowAligned, AlignedRecordDuplicateMeasurementWritesOneRow, AlignedRegisterAfterWriteIsRejected and RestorableTsFileIOWriterTest.AlignedTimeseriesRecoverAndWriteNullValue. --- cpp/src/writer/tsfile_writer.cc | 213 +++++++--- cpp/src/writer/tsfile_writer.h | 8 + cpp/src/writer/value_chunk_writer.h | 35 ++ cpp/src/writer/value_page_writer.h | 20 + .../file/restorable_tsfile_io_writer_test.cc | 126 ++++++ cpp/test/writer/tsfile_writer_test.cc | 385 ++++++++++++++++++ 6 files changed, 740 insertions(+), 47 deletions(-) diff --git a/cpp/src/writer/tsfile_writer.cc b/cpp/src/writer/tsfile_writer.cc index 41a485c00..6c0ebacf6 100644 --- a/cpp/src/writer/tsfile_writer.cc +++ b/cpp/src/writer/tsfile_writer.cc @@ -21,6 +21,8 @@ #include #include +#include +#include #include "chunk_writer.h" #include "common/config/config.h" @@ -334,8 +336,30 @@ int TsFileWriter::register_timeseries(const std::string& device_path, std::make_shared(device_path); DeviceSchemasMapIter device_iter = schemas_.find(device_id); if (device_iter != schemas_.end()) { - MeasurementSchemaMap& msm = - device_iter->second->measurement_schema_map_; + MeasurementSchemaGroup* device_schema = device_iter->second; + MeasurementSchemaMap& msm = device_schema->measurement_schema_map_; + if (msm.find(measurement_schema->measurement_name_) != msm.end()) { + return E_ALREADY_EXIST; + } + if (device_schema->is_aligned_ && + device_schema->time_chunk_writer_ != nullptr && + device_schema->time_chunk_writer_->hasData()) { + // A column added now would have to be padded with the rows (and + // pages) that were already written, which the current page writer + // cannot express. Java does not allow an aligned device to be + // expanded at all; refuse loudly instead of silently writing a + // chunk group whose value column row counts diverge from the time + // column. + return E_INVALID_ARG; + } + // Aligned devices advance every registered measurement on every row, + // so the value chunk writer has to exist before the first row. + if (device_schema->is_aligned_) { + int ret = ensure_aligned_value_chunk_writer(measurement_schema); + if (RET_FAIL(ret)) { + return ret; + } + } MeasurementSchemaMapInsertResult ins_res = msm.insert(std::make_pair( measurement_schema->measurement_name_, measurement_schema)); if (UNLIKELY(!ins_res.second)) { @@ -344,6 +368,13 @@ int TsFileWriter::register_timeseries(const std::string& device_path, } else { MeasurementSchemaGroup* ms_group = new MeasurementSchemaGroup; ms_group->is_aligned_ = is_aligned; + if (is_aligned) { + int ret = ensure_aligned_value_chunk_writer(measurement_schema); + if (RET_FAIL(ret)) { + delete ms_group; + return ret; + } + } ms_group->measurement_schema_map_.insert(std::make_pair( measurement_schema->measurement_name_, measurement_schema)); schemas_.insert(std::make_pair(device_id, ms_group)); @@ -516,6 +547,26 @@ int TsFileWriter::do_check_schema( return ret; } +int TsFileWriter::ensure_aligned_value_chunk_writer( + MeasurementSchema* measurement_schema) { + if (measurement_schema->value_chunk_writer_ != nullptr) { + return E_OK; + } + ValueChunkWriter* value_chunk_writer = new ValueChunkWriter; + if (IS_NULL(value_chunk_writer)) { + return E_OOM; + } + int ret = value_chunk_writer->init( + measurement_schema->measurement_name_, measurement_schema->data_type_, + measurement_schema->encoding_, measurement_schema->compression_type_); + if (RET_FAIL(ret)) { + delete value_chunk_writer; + return (ret == E_OOM) ? ret : common::E_INVALID_ARG; + } + measurement_schema->value_chunk_writer_ = value_chunk_writer; + return E_OK; +} + template int TsFileWriter::do_check_schema_aligned( std::shared_ptr device_id, @@ -548,28 +599,11 @@ int TsFileWriter::do_check_schema_aligned( // Here we may check data_type against ms_iter. But in Java // libtsfile, no check here. MeasurementSchema* ms = ms_iter->second; - if (IS_NULL(ms->value_chunk_writer_)) { - ms->value_chunk_writer_ = new ValueChunkWriter; - ret = ms->value_chunk_writer_->init( - ms->measurement_name_, ms->data_type_, ms->encoding_, - ms->compression_type_); - if (IS_SUCC(ret)) { - value_chunk_writers.push_back(ms->value_chunk_writer_); - } else { - value_chunk_writers.push_back(NULL); - for (size_t chunk_writer_idx = 0; - chunk_writer_idx < value_chunk_writers.size(); - chunk_writer_idx++) { - if (!value_chunk_writers[chunk_writer_idx]) { - delete value_chunk_writers[chunk_writer_idx]; - } - } - ret = common::E_INVALID_ARG; - return ret; - } - } else { - value_chunk_writers.push_back(ms->value_chunk_writer_); + if (RET_FAIL(ensure_aligned_value_chunk_writer(ms))) { + value_chunk_writers.push_back(NULL); + return common::E_INVALID_ARG; } + value_chunk_writers.push_back(ms->value_chunk_writer_); data_types.push_back(ms->data_type_); } } @@ -808,15 +842,50 @@ int TsFileWriter::write_record_aligned(const TsRecord& record) { if (value_chunk_writers.size() != record.points_.size()) { return E_INVALID_ARG; } + DeviceSchemasMapIter dev_it = schemas_.find(device_id); + if (UNLIKELY(dev_it == schemas_.end()) || IS_NULL(dev_it->second)) { + return E_DEVICE_NOT_EXIST; + } + MeasurementSchemaGroup* device_schema = dev_it->second; + // Index of the point that carries each measurement of this row. A + // duplicate point for the same measurement keeps the last one, so that + // every value column advances exactly one row per timestamp. + std::map row_point_index; + for (uint32_t c = 0; c < record.points_.size(); c++) { + row_point_index[record.points_[c].measurement_name_] = c; + } + // A record may only carry a subset of the device's measurements. The + // measurements it does not mention still have to advance one row (as a + // NULL) so that every value column's not-null bitmap stays aligned with + // the time column; otherwise the values of that column end up paired with + // the wrong timestamps on read. Java behaves the same way + // (AlignedChunkGroupWriterImpl#write -> writeEmptyDataInOneRow). + SimpleVector absent_schemas; + SimpleVector row_writers; + for (uint32_t c = 0; c < value_chunk_writers.size(); c++) { + if (!IS_NULL(value_chunk_writers[c])) { + row_writers.push_back(value_chunk_writers[c]); + } + } + MeasurementSchemaMap& msm = device_schema->measurement_schema_map_; + for (MeasurementSchemaMapIter ms_iter = msm.begin(); ms_iter != msm.end(); + ms_iter++) { + if (row_point_index.find(ms_iter->first) != row_point_index.end()) { + continue; + } + MeasurementSchema* ms = ms_iter->second; + if (RET_FAIL(ensure_aligned_value_chunk_writer(ms))) { + return ret; + } + absent_schemas.push_back(ms); + row_writers.push_back(ms->value_chunk_writer_); + } // Snapshot page counters before the write so we can detect any column // that crossed a page boundary and seal the rest in lockstep. int32_t time_pages_before = time_chunk_writer->num_of_pages(); - std::vector value_pages_before(value_chunk_writers.size(), 0); - for (uint32_t c = 0; c < value_chunk_writers.size(); c++) { - ValueChunkWriter* value_chunk_writer = value_chunk_writers[c]; - if (!IS_NULL(value_chunk_writer)) { - value_pages_before[c] = value_chunk_writer->num_of_pages(); - } + std::vector value_pages_before(row_writers.size(), 0); + for (uint32_t c = 0; c < row_writers.size(); c++) { + value_pages_before[c] = row_writers[c]->num_of_pages(); } // Time first: a rejected timestamp (E_OUT_OF_ORDER, OOM, etc.) must // not silently advance the value writers — that would leave the time @@ -829,8 +898,14 @@ int TsFileWriter::write_record_aligned(const TsRecord& record) { if (IS_NULL(value_chunk_writer)) { continue; } + const DataPoint& point = record.points_[c]; + if (row_point_index[point.measurement_name_] != c) { + // Duplicate point for this measurement: skipped, the last one is + // written below (one row per column per timestamp). + continue; + } if (RET_FAIL(write_point_aligned(value_chunk_writer, record.timestamp_, - data_types[c], record.points_[c]))) { + data_types[c], point))) { // Time wrote the row but at least one value column failed // mid-record; the per-column row counts no longer agree. // Mark the writer unrecoverable so flush/close refuses to @@ -839,8 +914,15 @@ int TsFileWriter::write_record_aligned(const TsRecord& record) { return ret; } } + for (uint32_t c = 0; c < absent_schemas.size(); c++) { + if (RET_FAIL( + absent_schemas[c]->value_chunk_writer_->write_null_batch(1))) { + unrecoverable_ = true; + return ret; + } + } if (RET_FAIL(maybe_seal_aligned_pages_together( - time_chunk_writer, value_chunk_writers, time_pages_before, + time_chunk_writer, row_writers, time_pages_before, value_pages_before))) { unrecoverable_ = true; return ret; @@ -984,15 +1066,50 @@ int TsFileWriter::write_tablet_aligned(const Tablet& tablet) { return E_TYPE_NOT_MATCH; } } + DeviceSchemasMapIter dev_it = schemas_.find(device_id); + if (UNLIKELY(dev_it == schemas_.end()) || IS_NULL(dev_it->second)) { + return E_DEVICE_NOT_EXIST; + } + MeasurementSchemaGroup* device_schema = dev_it->second; + // Measurements of this aligned device that the tablet does not carry still + // have to advance one row per tablet row (as NULLs), otherwise their + // not-null bitmap would start at the wrong row index and their values + // would be paired with the wrong timestamps on read. Java fills those + // rows in AlignedChunkGroupWriterImpl#write(Tablet). + SimpleVector absent_writers; + std::set tablet_measurements; + for (size_t c = 0; c < tablet.get_column_count(); c++) { + tablet_measurements.insert(tablet.schema_vec_->at(c).measurement_name_); + } + MeasurementSchemaMap& msm = device_schema->measurement_schema_map_; + for (MeasurementSchemaMapIter ms_iter = msm.begin(); ms_iter != msm.end(); + ms_iter++) { + if (tablet_measurements.find(ms_iter->first) != + tablet_measurements.end()) { + continue; + } + if (RET_FAIL(ensure_aligned_value_chunk_writer(ms_iter->second))) { + return ret; + } + absent_writers.push_back(ms_iter->second->value_chunk_writer_); + } + // Every column that takes part in this batch: the tablet's own columns + // plus the ones that only advance with NULL rows. + SimpleVector row_writers; + for (uint32_t c = 0; c < value_chunk_writers.size(); c++) { + if (!IS_NULL(value_chunk_writers[c])) { + row_writers.push_back(value_chunk_writers[c]); + } + } + for (uint32_t c = 0; c < absent_writers.size(); c++) { + row_writers.push_back(absent_writers[c]); + } // Snapshot page counters before the batch so we can detect any column // that crossed a page boundary mid-tablet and seal the rest in lockstep. int32_t time_pages_before = time_chunk_writer->num_of_pages(); - std::vector value_pages_before(value_chunk_writers.size(), 0); - for (uint32_t c = 0; c < value_chunk_writers.size(); c++) { - ValueChunkWriter* value_chunk_writer = value_chunk_writers[c]; - if (!IS_NULL(value_chunk_writer)) { - value_pages_before[c] = value_chunk_writer->num_of_pages(); - } + std::vector value_pages_before(row_writers.size(), 0); + for (uint32_t c = 0; c < row_writers.size(); c++) { + value_pages_before[c] = row_writers[c]->num_of_pages(); } // Suppress memory-driven page sealing on every column for the duration of // the batch. The count-driven seals inside write_batch still fire at the @@ -1004,18 +1121,13 @@ int TsFileWriter::write_tablet_aligned(const Tablet& tablet) { // (e.g. when a sealed value column ended a page that the time column did // not). time_chunk_writer->set_enable_page_seal_if_full(false); - for (uint32_t c = 0; c < value_chunk_writers.size(); c++) { - ValueChunkWriter* value_chunk_writer = value_chunk_writers[c]; - if (!IS_NULL(value_chunk_writer)) { - value_chunk_writer->set_enable_page_seal_if_full(false); - } + for (uint32_t c = 0; c < row_writers.size(); c++) { + row_writers[c]->set_enable_page_seal_if_full(false); } auto restore_seal = [&]() { time_chunk_writer->set_enable_page_seal_if_full(true); - for (uint32_t k = 0; k < value_chunk_writers.size(); k++) { - if (!IS_NULL(value_chunk_writers[k])) { - value_chunk_writers[k]->set_enable_page_seal_if_full(true); - } + for (uint32_t k = 0; k < row_writers.size(); k++) { + row_writers[k]->set_enable_page_seal_if_full(true); } }; // Any failure (out-of-order timestamps, OOM, etc.) must abort before we @@ -1043,9 +1155,16 @@ int TsFileWriter::write_tablet_aligned(const Tablet& tablet) { return ret; } } + for (uint32_t c = 0; c < absent_writers.size(); c++) { + if (RET_FAIL(absent_writers[c]->write_null_batch(total_rows))) { + restore_seal(); + unrecoverable_ = true; + return ret; + } + } restore_seal(); if (RET_FAIL(maybe_seal_aligned_pages_together( - time_chunk_writer, value_chunk_writers, time_pages_before, + time_chunk_writer, row_writers, time_pages_before, value_pages_before))) { unrecoverable_ = true; return ret; diff --git a/cpp/src/writer/tsfile_writer.h b/cpp/src/writer/tsfile_writer.h index 55e9e7f3a..0737eb065 100644 --- a/cpp/src/writer/tsfile_writer.h +++ b/cpp/src/writer/tsfile_writer.h @@ -128,6 +128,14 @@ class TsFileWriter { int write_point_aligned(ValueChunkWriter* value_chunk_writer, int64_t timestamp, common::TSDataType data_type, const DataPoint& point); + /* + * Create (once) the value chunk writer that carries one measurement of an + * aligned device. Aligned devices keep every registered measurement in + * lock-step with the time column, so the writer has to exist before the + * first row of that device is written. + */ + int ensure_aligned_value_chunk_writer( + storage::MeasurementSchema* measurement_schema); int maybe_seal_aligned_pages_together( TimeChunkWriter* time_chunk_writer, common::SimpleVector& value_chunk_writers, diff --git a/cpp/src/writer/value_chunk_writer.h b/cpp/src/writer/value_chunk_writer.h index cd7c75a54..88871c1a9 100644 --- a/cpp/src/writer/value_chunk_writer.h +++ b/cpp/src/writer/value_chunk_writer.h @@ -174,6 +174,41 @@ class ValueChunkWriter { return ret; } + /** + * Advance this column by `count` all-NULL rows. + * + * Aligned chunk groups keep every value column row-aligned with the time + * column, so a column whose row has no value still has to consume one row + * (with a null bit) for that timestamp. Page boundaries follow the same + * `page_writer_max_point_num_` rule as write_batch() so the page lists of + * the time column and of every value column stay in step. + */ + int write_null_batch(uint32_t count) { + int ret = common::E_OK; + uint32_t offset = 0; + const uint32_t page_cap = + common::g_config_value_.page_writer_max_point_num_; + while (offset < count) { + uint32_t cur_points = value_page_writer_.get_point_numer(); + if (cur_points >= page_cap) { + if (RET_FAIL(seal_cur_page(false))) { + return ret; + } + cur_points = 0; + } + uint32_t batch_size = + std::min(count - offset, page_cap - cur_points); + if (RET_FAIL(value_page_writer_.write_null_rows(batch_size))) { + return ret; + } + offset += batch_size; + if (RET_FAIL(seal_cur_page_if_full())) { + return ret; + } + } + return ret; + } + int end_encode_chunk(); common::ByteStream& get_chunk_data() { return chunk_data_; } Statistic* get_chunk_statistic() { return chunk_statistic_; } diff --git a/cpp/src/writer/value_page_writer.h b/cpp/src/writer/value_page_writer.h index 92c39b9b2..61bf59564 100644 --- a/cpp/src/writer/value_page_writer.h +++ b/cpp/src/writer/value_page_writer.h @@ -335,6 +335,26 @@ class ValuePageWriter { value_out_stream_.allocated_bytes()) + value_encoder_->get_max_byte_size(); } + /** + * Append `count` all-NULL rows to the current page: one zero bit per row + * in the not-null bitmap, no value bytes, no statistic update. + * + * Used to keep a value column row-aligned with the time column when a row + * carries no value for that column. The null bit has to be recorded for + * that row, otherwise the column's bitmap would start at the wrong row + * index and every later value would be paired with the wrong timestamp on + * read. + */ + int write_null_rows(uint32_t count) { + for (uint32_t i = 0; i < count; i++) { + if ((size_ / 8) + 1 > col_notnull_bitmap_.size()) { + col_notnull_bitmap_.push_back(0); + } + size_++; + } + return common::E_OK; + } + int write_to_chunk(common::ByteStream& pages_data, bool write_header, bool write_statistic, bool write_data_to_chunk_data); FORCE_INLINE common::ByteStream& get_col_notnull_bitmap_data() { diff --git a/cpp/test/file/restorable_tsfile_io_writer_test.cc b/cpp/test/file/restorable_tsfile_io_writer_test.cc index c60a855c5..137f24f08 100644 --- a/cpp/test/file/restorable_tsfile_io_writer_test.cc +++ b/cpp/test/file/restorable_tsfile_io_writer_test.cc @@ -1061,3 +1061,129 @@ TEST_F(RestorableTsFileIOWriterTest, RecoveryAlignedSparseStatRespectsBitmap) { } EXPECT_TRUE(found_value_chunk); } + +// Sparse aligned records (a row only carries a subset of the device's +// measurements, without explicit NULL DataPoints) must survive a crash + +// recovery: the recovered chunk statistics have to describe the rows that +// really carry a value, and continued sparse writes have to stay row-aligned +// with the time column. +TEST_F(RestorableTsFileIOWriterTest, + AlignedTimeseriesRecoverAndWriteNullValue) { + using namespace std; + const string device = "d1"; + // even rows carry s1..s3, odd rows carry s4..s6 + vector even_names = {"s1", "s2", "s3"}; + vector odd_names = {"s4", "s5", "s6"}; + { + TsFileWriter tw; + ASSERT_EQ(tw.open(file_name_, GetWriteCreateFlags(), 0666), E_OK); + std::vector schemas; + schemas.push_back(new MeasurementSchema("s1", BOOLEAN)); + schemas.push_back(new MeasurementSchema("s2", INT32)); + schemas.push_back(new MeasurementSchema("s3", TEXT)); + schemas.push_back(new MeasurementSchema("s4", INT64)); + schemas.push_back(new MeasurementSchema("s5", FLOAT)); + schemas.push_back(new MeasurementSchema("s6", STRING)); + ASSERT_EQ(tw.register_aligned_timeseries(device, schemas), E_OK); + for (int i = 0; i < 10; i++) { + TsRecord record(i, device); + if (i % 2 == 0) { + record.add_point(even_names[0], true); + record.add_point(even_names[1], static_cast(i)); + record.add_point(even_names[2], "even"); + } else { + record.add_point(odd_names[0], static_cast(i)); + record.add_point(odd_names[1], static_cast(i)); + record.add_point(odd_names[2], "odd"); + } + ASSERT_EQ(tw.write_record_aligned(record), E_OK); + } + ASSERT_EQ(tw.flush(), E_OK); + ASSERT_EQ(tw.close(), E_OK); + } + + CorruptCurrentFileTail(3); + + RestorableTsFileIOWriter rw; + ASSERT_EQ(rw.open(file_name_, true), E_OK); + ASSERT_TRUE(rw.can_write()); + { + // Keep writing sparse rows after the recovery point. + TsFileTreeWriter tw2(&rw); + for (int i = 10; i < 20; i++) { + TsRecord record(i, device); + if (i % 2 == 0) { + record.add_point(even_names[0], true); + record.add_point(even_names[1], static_cast(i)); + record.add_point(even_names[2], "even"); + } else { + record.add_point(odd_names[0], static_cast(i)); + record.add_point(odd_names[1], static_cast(i)); + record.add_point(odd_names[2], "odd"); + } + ASSERT_EQ(tw2.write(record), E_OK); + } + ASSERT_EQ(tw2.flush(), E_OK); + ASSERT_EQ(tw2.close(), E_OK); + } + + TsFileTreeReader reader; + ASSERT_EQ(reader.open(file_name_), E_OK); + DeviceTimeseriesMetadataMap metadata = reader.get_timeseries_metadata(); + std::map meta_by_name; + for (auto& entry : metadata) { + for (auto& ts_idx : entry.second) { + meta_by_name[ts_idx->get_measurement_name().to_std_string()] = + ts_idx.get(); + } + } + ASSERT_EQ(meta_by_name.size(), 6u); + // Columns carried by the even rows: 5 points before the crash + 5 after, + // spanning timestamp 0..18. + for (auto& name : even_names) { + EXPECT_EQ(meta_by_name[name]->get_statistic()->count_, 10); + EXPECT_EQ(meta_by_name[name]->get_statistic()->start_time_, 0); + EXPECT_EQ(meta_by_name[name]->get_statistic()->end_time_, 18); + } + // Columns carried by the odd rows: first value at timestamp 1, last at 19. + for (auto& name : odd_names) { + EXPECT_EQ(meta_by_name[name]->get_statistic()->count_, 10); + EXPECT_EQ(meta_by_name[name]->get_statistic()->start_time_, 1); + EXPECT_EQ(meta_by_name[name]->get_statistic()->end_time_, 19); + } + + vector measurement_names = {"s1", "s2", "s3", "s4", "s5", "s6"}; + ASSERT_EQ(CountTreeReaderRows(reader, measurement_names), 20); + + ResultSet* result_set = nullptr; + vector device_ids = {device}; + ASSERT_EQ(reader.query(device_ids, measurement_names, 0, 100, result_set), + E_OK); + auto it = result_set->iterator(); + int row = 0; + while (it.hasNext()) { + RowRecord* rec = it.next(); + ASSERT_NE(rec, nullptr); + EXPECT_EQ(rec->get_timestamp(), row); + for (int c = 0; c < 3; c++) { + Field* field = rec->get_field(c + 1); + if (row % 2 == 0) { + EXPECT_NE(field->type_, common::NULL_TYPE); + } else { + EXPECT_EQ(field->type_, common::NULL_TYPE); + } + } + for (int c = 3; c < 6; c++) { + Field* field = rec->get_field(c + 1); + if (row % 2 == 0) { + EXPECT_EQ(field->type_, common::NULL_TYPE); + } else { + EXPECT_NE(field->type_, common::NULL_TYPE); + } + } + row++; + } + EXPECT_EQ(row, 20); + reader.destroy_query_data_set(result_set); + reader.close(); +} diff --git a/cpp/test/writer/tsfile_writer_test.cc b/cpp/test/writer/tsfile_writer_test.cc index 3b9dae92a..ef9527ff9 100644 --- a/cpp/test/writer/tsfile_writer_test.cc +++ b/cpp/test/writer/tsfile_writer_test.cc @@ -1737,3 +1737,388 @@ TEST_F(TsFileWriterTest, WriterReuseAfterDestroyProducesValidSecondFile) { delete wf; remove(second_path.c_str()); } + +// --------------------------------------------------------------------------- +// Aligned writes: a row that does not carry every measurement of the device +// must not shift the value columns. +// +// Records and tablets may legitimately carry only a subset of an aligned +// device's measurements. A measurement missing from the row still has to +// consume one row (as NULL) so that every value column's not-null bitmap stays +// aligned with the time column; otherwise the values of that column get paired +// with the earliest timestamps of the page on read (Java behaves the same way: +// AlignedChunkGroupWriterImpl#write -> writeEmptyDataInOneRow). +// --------------------------------------------------------------------------- + +TEST_F(TsFileWriterTest, AlignedRecordMissingMeasurementsStayRowAligned) { + std::string device_name = "device_missing_m"; + std::vector mnames = {"s0", "s1", "s2"}; + std::vector schemas; + for (auto& n : mnames) { + schemas.push_back(new MeasurementSchema(n, INT64, PLAIN, UNCOMPRESSED)); + } + tsfile_writer_->register_aligned_timeseries(device_name, schemas); + + const int row_num = 10; + for (int i = 0; i < row_num; i++) { + TsRecord record(1622505600000 + i, device_name); + if (i % 2 == 0) { + // Only s0 is carried by even rows, only s1 by odd rows; s2 is + // never written at all. + record.add_point(mnames[0], static_cast(100 + i)); + } else { + record.add_point(mnames[1], static_cast(200 + i)); + } + ASSERT_EQ(tsfile_writer_->write_record_aligned(record), E_OK); + } + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); + + // Statistics must describe the rows that really carry a value. + std::vector> devices = { + std::make_shared(device_name)}; + TsFileReader meta_reader; + ASSERT_EQ(meta_reader.open(file_name_), E_OK); + auto meta_map = meta_reader.get_timeseries_metadata(devices); + std::map meta_by_name; + for (auto& ts_idx : meta_map.at(devices[0])) { + meta_by_name[ts_idx->get_measurement_name().to_std_string()] = + ts_idx.get(); + } + // s2 is registered but never carries a value: it advances with NULL rows + // (like an all-null column of an aligned tablet), so it is present with + // count 0. + ASSERT_EQ(meta_by_name.size(), 3u); + EXPECT_EQ(meta_by_name[mnames[0]]->get_statistic()->count_, 5); + EXPECT_EQ(meta_by_name[mnames[0]]->get_statistic()->start_time_, + 1622505600000); + EXPECT_EQ(meta_by_name[mnames[0]]->get_statistic()->end_time_, + 1622505600008); + EXPECT_EQ(meta_by_name[mnames[1]]->get_statistic()->count_, 5); + EXPECT_EQ(meta_by_name[mnames[1]]->get_statistic()->start_time_, + 1622505600001); + EXPECT_EQ(meta_by_name[mnames[1]]->get_statistic()->end_time_, + 1622505600009); + EXPECT_EQ(meta_by_name[mnames[2]]->get_statistic()->count_, 0); + ASSERT_EQ(meta_reader.close(), E_OK); + + std::vector select_list; + for (auto& n : mnames) { + select_list.emplace_back(device_name, n); + } + storage::QueryExpression* qe = + storage::QueryExpression::create(select_list, nullptr); + storage::TsFileReader reader; + ASSERT_EQ(reader.open(file_name_), E_OK); + storage::ResultSet* tmp_qds = nullptr; + ASSERT_EQ(reader.query(qe, tmp_qds), E_OK); + auto* qds = (QDSWithoutTimeGenerator*)tmp_qds; + + bool has_next = false; + int64_t cur_row = 0; + while (IS_SUCC(qds->next(has_next)) && has_next) { + auto* rec = qds->get_row_record(); + ASSERT_NE(rec, nullptr); + EXPECT_EQ(rec->get_timestamp(), 1622505600000 + cur_row); + const std::string s0 = field_to_string(rec->get_field(1)); + const std::string s1 = field_to_string(rec->get_field(2)); + const std::string s2 = field_to_string(rec->get_field(3)); + if (cur_row % 2 == 0) { + EXPECT_EQ(s0, std::to_string(100 + cur_row)); + EXPECT_EQ(s1, "NULL"); + } else { + EXPECT_EQ(s0, "NULL"); + EXPECT_EQ(s1, std::to_string(200 + cur_row)); + } + EXPECT_EQ(s2, "NULL"); + cur_row++; + } + EXPECT_EQ(cur_row, row_num); + reader.destroy_query_data_set(qds); + ASSERT_EQ(reader.close(), E_OK); +} + +TEST_F(TsFileWriterTest, AlignedTabletMissingColumnStaysRowAligned) { + std::string device_name = "device_tablet_missing"; + std::vector schema_vec; + schema_vec.emplace_back("s0", INT64, PLAIN, UNCOMPRESSED); + schema_vec.emplace_back("s1", INT64, PLAIN, UNCOMPRESSED); + { + std::vector reg; + for (auto& s : schema_vec) { + reg.push_back(new MeasurementSchema(s)); + } + tsfile_writer_->register_aligned_timeseries(device_name, reg); + } + { + // First tablet only carries s0: s1 must still advance with NULLs. + std::vector cols; + cols.push_back(schema_vec[0]); + Tablet tablet(device_name, + std::make_shared>(cols), + 5); + for (int i = 0; i < 5; i++) { + tablet.add_timestamp(i, 1000 + i); + tablet.add_value(i, 0u, static_cast(100 + i)); + } + ASSERT_EQ(tsfile_writer_->write_tablet_aligned(tablet), E_OK); + } + { + Tablet tablet( + device_name, + std::make_shared>(schema_vec), 5); + for (int i = 0; i < 5; i++) { + tablet.add_timestamp(i, 1005 + i); + tablet.add_value(i, 0u, static_cast(105 + i)); + tablet.add_value(i, 1u, static_cast(200 + i)); + } + ASSERT_EQ(tsfile_writer_->write_tablet_aligned(tablet), E_OK); + } + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); + + std::string s0_name("s0"), s1_name("s1"); + std::vector select_list; + select_list.emplace_back(device_name, s0_name); + select_list.emplace_back(device_name, s1_name); + storage::QueryExpression* qe = + storage::QueryExpression::create(select_list, nullptr); + storage::TsFileReader reader; + ASSERT_EQ(reader.open(file_name_), E_OK); + storage::ResultSet* tmp_qds = nullptr; + ASSERT_EQ(reader.query(qe, tmp_qds), E_OK); + auto* qds = (QDSWithoutTimeGenerator*)tmp_qds; + + bool has_next = false; + int64_t cur_row = 0; + while (IS_SUCC(qds->next(has_next)) && has_next) { + auto* rec = qds->get_row_record(); + ASSERT_NE(rec, nullptr); + EXPECT_EQ(rec->get_timestamp(), 1000 + cur_row); + EXPECT_EQ(field_to_string(rec->get_field(1)), + std::to_string(100 + cur_row)); + if (cur_row < 5) { + // The rows written before s1 showed up are NULL for s1. + EXPECT_EQ(field_to_string(rec->get_field(2)), "NULL"); + } else { + EXPECT_EQ(field_to_string(rec->get_field(2)), + std::to_string(200 + cur_row - 5)); + } + cur_row++; + } + EXPECT_EQ(cur_row, 10); + reader.destroy_query_data_set(qds); + ASSERT_EQ(reader.close(), E_OK); +} + +// A record that repeats the same measurement must still advance that column a +// single time, otherwise the column would run ahead of the time column. +TEST_F(TsFileWriterTest, AlignedRecordDuplicateMeasurementWritesOneRow) { + std::string device_name = "device_dup_m"; + std::vector schemas; + schemas.push_back(new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED)); + schemas.push_back(new MeasurementSchema("s1", INT64, PLAIN, UNCOMPRESSED)); + tsfile_writer_->register_aligned_timeseries(device_name, schemas); + + TsRecord record(7, device_name); + record.add_point("s0", static_cast(1)); + record.add_point("s0", static_cast(2)); + record.add_point("s1", static_cast(3)); + ASSERT_EQ(tsfile_writer_->write_record_aligned(record), E_OK); + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); + + std::string s0_name("s0"), s1_name("s1"); + std::vector select_list; + select_list.emplace_back(device_name, s0_name); + select_list.emplace_back(device_name, s1_name); + storage::QueryExpression* qe = + storage::QueryExpression::create(select_list, nullptr); + storage::TsFileReader reader; + ASSERT_EQ(reader.open(file_name_), E_OK); + storage::ResultSet* tmp_qds = nullptr; + ASSERT_EQ(reader.query(qe, tmp_qds), E_OK); + auto* qds = (QDSWithoutTimeGenerator*)tmp_qds; + + bool has_next = false; + int rows = 0; + while (IS_SUCC(qds->next(has_next)) && has_next) { + auto* rec = qds->get_row_record(); + ASSERT_NE(rec, nullptr); + EXPECT_EQ(rec->get_timestamp(), 7); + // The last point of the duplicated measurement wins. + EXPECT_EQ(field_to_string(rec->get_field(1)), "2"); + EXPECT_EQ(field_to_string(rec->get_field(2)), "3"); + rows++; + } + EXPECT_EQ(rows, 1); + reader.destroy_query_data_set(qds); + ASSERT_EQ(reader.close(), E_OK); +} + +// New columns have to be padded with the rows (and pages) that were already +// written, which the aligned writer cannot express: reject the registration +// instead of writing a chunk group whose value column is shifted. +TEST_F(TsFileWriterTest, AlignedRegisterAfterWriteIsRejected) { + std::string device_name = "device_late_m"; + std::vector schemas; + schemas.push_back(new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED)); + tsfile_writer_->register_aligned_timeseries(device_name, schemas); + + TsRecord record(1, device_name); + record.add_point("s0", static_cast(1)); + ASSERT_EQ(tsfile_writer_->write_record_aligned(record), E_OK); + + std::vector extra; + extra.push_back(new MeasurementSchema("s1", INT64, PLAIN, UNCOMPRESSED)); + EXPECT_EQ(tsfile_writer_->register_aligned_timeseries(device_name, extra), + E_INVALID_ARG); + + // A second registration of the same measurement is still reported as a + // duplicate, not as a late registration. + std::vector dup; + dup.push_back(new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED)); + EXPECT_EQ(tsfile_writer_->register_aligned_timeseries(device_name, dup), + E_ALREADY_EXIST); + + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); +} + +// Same as above, but with rows spread over several pages: the NULL padding of +// a missing measurement has to keep the page lists of every value column in +// step with the time column. +TEST_F(TsFileWriterTest, AlignedRecordMissingMeasurementsAcrossPages) { + uint32_t prev_pt = g_config_value_.page_writer_max_point_num_; + uint32_t prev_mem = g_config_value_.page_writer_max_memory_bytes_; + struct Guard { + uint32_t pt, mem; + ~Guard() { + g_config_value_.page_writer_max_point_num_ = pt; + g_config_value_.page_writer_max_memory_bytes_ = mem; + } + } guard{prev_pt, prev_mem}; + g_config_value_.page_writer_max_point_num_ = 7; + g_config_value_.page_writer_max_memory_bytes_ = 1024 * 1024; + + std::string device_name = "device_missing_pages"; + std::vector mnames = {"s0", "s1", "s2"}; + std::vector schemas; + for (auto& n : mnames) { + schemas.push_back(new MeasurementSchema(n, INT64, PLAIN, UNCOMPRESSED)); + } + tsfile_writer_->register_aligned_timeseries(device_name, schemas); + + const int row_num = 20; + for (int i = 0; i < row_num; i++) { + TsRecord record(1000 + i, device_name); + // s0: every row, s1: every third row, s2: only the last row. + record.add_point(mnames[0], static_cast(i)); + if (i % 3 == 0) { + record.add_point(mnames[1], static_cast(100 + i)); + } + if (i == row_num - 1) { + record.add_point(mnames[2], static_cast(999)); + } + ASSERT_EQ(tsfile_writer_->write_record_aligned(record), E_OK); + } + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); + + std::vector select_list; + for (auto& n : mnames) { + select_list.emplace_back(device_name, n); + } + storage::QueryExpression* qe = + storage::QueryExpression::create(select_list, nullptr); + storage::TsFileReader reader; + ASSERT_EQ(reader.open(file_name_), E_OK); + storage::ResultSet* tmp_qds = nullptr; + ASSERT_EQ(reader.query(qe, tmp_qds), E_OK); + auto* qds = (QDSWithoutTimeGenerator*)tmp_qds; + + bool has_next = false; + int64_t cur_row = 0; + while (IS_SUCC(qds->next(has_next)) && has_next) { + auto* rec = qds->get_row_record(); + ASSERT_NE(rec, nullptr); + EXPECT_EQ(rec->get_timestamp(), 1000 + cur_row); + EXPECT_EQ(field_to_string(rec->get_field(1)), std::to_string(cur_row)); + if (cur_row % 3 == 0) { + EXPECT_EQ(field_to_string(rec->get_field(2)), + std::to_string(100 + cur_row)); + } else { + EXPECT_EQ(field_to_string(rec->get_field(2)), "NULL"); + } + if (cur_row == row_num - 1) { + EXPECT_EQ(field_to_string(rec->get_field(3)), "999"); + } else { + EXPECT_EQ(field_to_string(rec->get_field(3)), "NULL"); + } + cur_row++; + } + EXPECT_EQ(cur_row, row_num); + reader.destroy_query_data_set(qds); + ASSERT_EQ(reader.close(), E_OK); +} + +// A value column is identified by measurement name, not by the position of the +// point inside the record: records may add their points in any order (and in a +// different order from row to row) without disturbing the row alignment. +TEST_F(TsFileWriterTest, AlignedRecordPointOrderDoesNotMatter) { + std::string device_name = "device_point_order"; + std::vector mnames = {"s0", "s1", "s2"}; + std::vector schemas; + for (auto& n : mnames) { + schemas.push_back(new MeasurementSchema(n, INT64, PLAIN, UNCOMPRESSED)); + } + tsfile_writer_->register_aligned_timeseries(device_name, schemas); + + const int row_num = 6; + for (int i = 0; i < row_num; i++) { + TsRecord record(2000 + i, device_name); + // Reverse order of the previous row, and drop s1 on odd rows. + for (int k = 2; k >= 0; k--) { + int idx = (i + k) % 3; + if (idx == 1 && i % 2 == 1) { + continue; + } + record.add_point(mnames[idx], static_cast(idx * 100 + i)); + } + ASSERT_EQ(tsfile_writer_->write_record_aligned(record), E_OK); + } + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); + + std::vector select_list; + for (auto& n : mnames) { + select_list.emplace_back(device_name, n); + } + storage::QueryExpression* qe = + storage::QueryExpression::create(select_list, nullptr); + storage::TsFileReader reader; + ASSERT_EQ(reader.open(file_name_), E_OK); + storage::ResultSet* tmp_qds = nullptr; + ASSERT_EQ(reader.query(qe, tmp_qds), E_OK); + auto* qds = (QDSWithoutTimeGenerator*)tmp_qds; + + bool has_next = false; + int64_t cur_row = 0; + while (IS_SUCC(qds->next(has_next)) && has_next) { + auto* rec = qds->get_row_record(); + ASSERT_NE(rec, nullptr); + EXPECT_EQ(rec->get_timestamp(), 2000 + cur_row); + for (int c = 0; c < 3; c++) { + if (c == 1 && cur_row % 2 == 1) { + EXPECT_EQ(field_to_string(rec->get_field(c + 1)), "NULL"); + } else { + EXPECT_EQ(field_to_string(rec->get_field(c + 1)), + std::to_string(c * 100 + cur_row)); + } + } + cur_row++; + } + EXPECT_EQ(cur_row, row_num); + reader.destroy_query_data_set(qds); + ASSERT_EQ(reader.close(), E_OK); +} From e5425e4fe1c3a8ed25df582e36a433783c32f90b Mon Sep 17 00:00:00 2001 From: ColinLee Date: Wed, 23 Sep 2026 10:58:04 +0800 Subject: [PATCH 2/4] test(cpp): release schemas whose aligned registration was rejected The writer only takes ownership of a MeasurementSchema when the registration succeeds, so the two schemas used to provoke E_INVALID_ARG / E_ALREADY_EXIST have to be released by the test. LeakSanitizer flagged them in the ASan jobs of PR #968 (208 bytes in 2 allocations). --- cpp/test/writer/tsfile_writer_test.cc | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/cpp/test/writer/tsfile_writer_test.cc b/cpp/test/writer/tsfile_writer_test.cc index ef9527ff9..b89ab7059 100644 --- a/cpp/test/writer/tsfile_writer_test.cc +++ b/cpp/test/writer/tsfile_writer_test.cc @@ -1969,17 +1969,23 @@ TEST_F(TsFileWriterTest, AlignedRegisterAfterWriteIsRejected) { record.add_point("s0", static_cast(1)); ASSERT_EQ(tsfile_writer_->write_record_aligned(record), E_OK); - std::vector extra; - extra.push_back(new MeasurementSchema("s1", INT64, PLAIN, UNCOMPRESSED)); + // The writer only takes ownership of a schema when the registration + // succeeds, so a rejected one has to be released by the caller. + MeasurementSchema* extra_schema = + new MeasurementSchema("s1", INT64, PLAIN, UNCOMPRESSED); + std::vector extra{extra_schema}; EXPECT_EQ(tsfile_writer_->register_aligned_timeseries(device_name, extra), E_INVALID_ARG); + delete extra_schema; // A second registration of the same measurement is still reported as a // duplicate, not as a late registration. - std::vector dup; - dup.push_back(new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED)); + MeasurementSchema* dup_schema = + new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED); + std::vector dup{dup_schema}; EXPECT_EQ(tsfile_writer_->register_aligned_timeseries(device_name, dup), E_ALREADY_EXIST); + delete dup_schema; ASSERT_EQ(tsfile_writer_->flush(), E_OK); ASSERT_EQ(tsfile_writer_->close(), E_OK); From d687fd8967578c15ec792b76db63818a0c6478d7 Mon Sep 17 00:00:00 2001 From: ColinLee Date: Sat, 10 Oct 2026 08:47:37 +0800 Subject: [PATCH 3/4] fix(cpp): seal full aligned tablet pages before record writes --- cpp/src/writer/tsfile_writer.cc | 25 ++++-- cpp/test/writer/tsfile_writer_test.cc | 117 ++++++++++++++++++++++++++ 2 files changed, 136 insertions(+), 6 deletions(-) diff --git a/cpp/src/writer/tsfile_writer.cc b/cpp/src/writer/tsfile_writer.cc index 6c0ebacf6..e72562136 100644 --- a/cpp/src/writer/tsfile_writer.cc +++ b/cpp/src/writer/tsfile_writer.cc @@ -305,7 +305,11 @@ int TsFileWriter::register_aligned_timeseries( MeasurementSchema* ms = new MeasurementSchema( measurement_schema.measurement_name_, measurement_schema.data_type_, measurement_schema.encoding_, measurement_schema.compression_type_); - return register_timeseries(device_id, ms, true); + int ret = register_timeseries(device_id, ms, true); + if (ret != E_OK) { + delete ms; + } + return ret; } int TsFileWriter::register_aligned_timeseries( @@ -326,7 +330,11 @@ int TsFileWriter::register_timeseries( MeasurementSchema* ms = new MeasurementSchema( measurement_schema.measurement_name_, measurement_schema.data_type_, measurement_schema.encoding_, measurement_schema.compression_type_); - return register_timeseries(device_id, ms, false); + int ret = register_timeseries(device_id, ms, false); + if (ret != E_OK) { + delete ms; + } + return ret; } int TsFileWriter::register_timeseries(const std::string& device_path, @@ -963,15 +971,20 @@ int TsFileWriter::write_point(ChunkWriter* chunk_writer, int64_t timestamp, } // After writing one record / batch to the time chunk and every value chunk, -// keep their page boundaries aligned: if any of them autosealed a page on -// memory pressure, seal the rest of the open pages too so an aligned reader -// can still pair position N across time + every value column. +// keep their page boundaries aligned: if any of them autosealed a page, or +// the batch left a full current page, seal the rest of the open pages too so +// an aligned reader can still pair position N across time + every value column. int TsFileWriter::maybe_seal_aligned_pages_together( TimeChunkWriter* time_chunk_writer, common::SimpleVector& value_chunk_writers, int32_t time_pages_before, const std::vector& value_pages_before) { + // Batch writes can leave the last page full without advancing the page + // count. Seal it before a later sparse record can split time and NULL + // values across different pages. bool should_seal_all = - time_chunk_writer->num_of_pages() > time_pages_before; + time_chunk_writer->num_of_pages() > time_pages_before || + time_chunk_writer->get_point_numer() >= + g_config_value_.page_writer_max_point_num_; for (uint32_t c = 0; c < value_chunk_writers.size() && !should_seal_all; c++) { ValueChunkWriter* value_chunk_writer = value_chunk_writers[c]; diff --git a/cpp/test/writer/tsfile_writer_test.cc b/cpp/test/writer/tsfile_writer_test.cc index b89ab7059..b2d1e4cce 100644 --- a/cpp/test/writer/tsfile_writer_test.cc +++ b/cpp/test/writer/tsfile_writer_test.cc @@ -24,6 +24,7 @@ #include #include #include +#include #ifdef _WIN32 #include @@ -1987,10 +1988,126 @@ TEST_F(TsFileWriterTest, AlignedRegisterAfterWriteIsRejected) { E_ALREADY_EXIST); delete dup_schema; + // The schema-reference overloads own their internal copies even when the + // registration is rejected. Exercise both wrappers under leak checking. + MeasurementSchema extra_value("s2", INT64, PLAIN, UNCOMPRESSED); + EXPECT_EQ( + tsfile_writer_->register_aligned_timeseries(device_name, extra_value), + E_INVALID_ARG); + EXPECT_EQ(tsfile_writer_->register_timeseries(device_name, extra_value), + E_INVALID_ARG); + MeasurementSchema duplicate_value("s0", INT64, PLAIN, UNCOMPRESSED); + EXPECT_EQ(tsfile_writer_->register_aligned_timeseries(device_name, + duplicate_value), + E_ALREADY_EXIST); + EXPECT_EQ(tsfile_writer_->register_timeseries(device_name, duplicate_value), + E_ALREADY_EXIST); + + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); +} + +class AlignedTabletRecordBoundaryTest + : public TsFileWriterTest, + public ::testing::WithParamInterface> {}; + +TEST_P(AlignedTabletRecordBoundaryTest, MissingMeasurementsStayAligned) { + const uint32_t tablet_rows = std::get<0>(GetParam()); + const bool missing_tablet_column = std::get<1>(GetParam()); + const int flush_mode = std::get<2>(GetParam()); + struct ConfigGuard { + uint32_t points, memory; + ~ConfigGuard() { + g_config_value_.page_writer_max_point_num_ = points; + g_config_value_.page_writer_max_memory_bytes_ = memory; + } + } guard{g_config_value_.page_writer_max_point_num_, + g_config_value_.page_writer_max_memory_bytes_}; + const uint32_t page_capacity = 7; + g_config_value_.page_writer_max_point_num_ = page_capacity; + g_config_value_.page_writer_max_memory_bytes_ = 1024 * 1024; + + std::string device_name = "device_tablet_record_boundary"; + std::vector schemas{ + new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED), + new MeasurementSchema("s1", INT64, PLAIN, UNCOMPRESSED)}; + ASSERT_EQ(tsfile_writer_->register_aligned_timeseries(device_name, schemas), + E_OK); + auto tablet_schema = std::make_shared>(); + tablet_schema->emplace_back("s0", INT64, PLAIN, UNCOMPRESSED); + if (!missing_tablet_column) { + tablet_schema->emplace_back("s1", INT64, PLAIN, UNCOMPRESSED); + } + Tablet tablet(device_name, tablet_schema, tablet_rows); + for (uint32_t row = 0; row < tablet_rows; row++) { + ASSERT_EQ(tablet.add_timestamp(row, row), E_OK); + ASSERT_EQ(tablet.add_value(row, 0u, int64_t(100 + row)), E_OK); + if (!missing_tablet_column) { + ASSERT_EQ(tablet.add_value(row, 1u, int64_t(200 + row)), E_OK); + } + } + ASSERT_EQ(tsfile_writer_->write_tablet_aligned(tablet), E_OK); + if (tablet_rows % page_capacity == 0) { + // A full last page must be sealed on every column before switching + // from batch writes to record writes, including all-NULL columns. + for (MeasurementSchema* schema : schemas) { + EXPECT_EQ(schema->value_chunk_writer_->get_point_numer(), 0u); + } + } + if (flush_mode == 1) { + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + } + + TsRecord sparse(tablet_rows, device_name); + sparse.add_point("s0", int64_t(100 + tablet_rows)); + ASSERT_EQ(tsfile_writer_->write_record_aligned(sparse), E_OK); + EXPECT_EQ(schemas[0]->value_chunk_writer_->num_of_pages(), + schemas[1]->value_chunk_writer_->num_of_pages()); + EXPECT_EQ(schemas[0]->value_chunk_writer_->get_point_numer(), + schemas[1]->value_chunk_writer_->get_point_numer()); + if (flush_mode == 2) { + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + } + + TsRecord full(tablet_rows + 1, device_name); + full.add_point("s0", int64_t(101 + tablet_rows)); + full.add_point("s1", int64_t(201 + tablet_rows)); + ASSERT_EQ(tsfile_writer_->write_record_aligned(full), E_OK); ASSERT_EQ(tsfile_writer_->flush(), E_OK); ASSERT_EQ(tsfile_writer_->close(), E_OK); + + std::string s0("s0"), s1("s1"); + std::vector paths{Path(device_name, s0), Path(device_name, s1)}; + TsFileReader reader; + ASSERT_EQ(reader.open(file_name_), E_OK); + ResultSet* result = nullptr; + ASSERT_EQ(reader.query(QueryExpression::create(paths, nullptr), result), + E_OK); + bool has_next = false; + uint32_t rows = 0; + int ret = E_OK; + while ((ret = result->next(has_next)) == E_OK && has_next) { + RowRecord* row = result->get_row_record(); + EXPECT_EQ(row->get_timestamp(), rows); + EXPECT_EQ(field_to_string(row->get_field(1)), + std::to_string(100 + rows)); + const bool is_null = rows == tablet_rows || + (missing_tablet_column && rows < tablet_rows); + EXPECT_EQ(field_to_string(row->get_field(2)), + is_null ? "NULL" : std::to_string(200 + rows)); + rows++; + } + EXPECT_EQ(ret, E_OK); + EXPECT_EQ(rows, tablet_rows + 2); + reader.destroy_query_data_set(result); + ASSERT_EQ(reader.close(), E_OK); } +INSTANTIATE_TEST_SUITE_P(PageBoundaries, AlignedTabletRecordBoundaryTest, + ::testing::Combine(::testing::Values(6u, 7u, 8u, 14u), + ::testing::Bool(), + ::testing::Values(0, 1, 2))); + // Same as above, but with rows spread over several pages: the NULL padding of // a missing measurement has to keep the page lists of every value column in // step with the time column. From 3a63e094ce32791125aa70ec5099f2f2a3d306c8 Mon Sep 17 00:00:00 2001 From: ColinLee Date: Sat, 10 Oct 2026 09:54:49 +0800 Subject: [PATCH 4/4] fix(cpp): freeze aligned measurements at device registration --- cpp/src/writer/tsfile_tree_writer.h | 8 +- cpp/src/writer/tsfile_writer.cc | 75 +++---- cpp/src/writer/tsfile_writer.h | 9 +- .../file/restorable_tsfile_io_writer_test.cc | 78 +++++++ cpp/test/reader/tsfile_reader_test.cc | 14 +- cpp/test/writer/tsfile_writer_test.cc | 198 +++++++++++++++--- 6 files changed, 302 insertions(+), 80 deletions(-) diff --git a/cpp/src/writer/tsfile_tree_writer.h b/cpp/src/writer/tsfile_tree_writer.h index 3c21c23da..f07046838 100644 --- a/cpp/src/writer/tsfile_tree_writer.h +++ b/cpp/src/writer/tsfile_tree_writer.h @@ -83,11 +83,15 @@ class TsFileTreeWriter { /** * Registers multiple aligned time series under the same device ID. + * The complete measurement list must be registered in one call. A device + * that is already registered or recovered cannot be registered again, + * including before writing or after flushing. * * @param device_id The ID or path of the device to which the aligned time * series belong. * @param schemas A vector of measurement schema pointers representing the - * aligned measurements. Must not be empty or contain null pointers. + * aligned measurements. Must not be empty or contain null pointers or + * duplicate names. Ownership transfers only on successful registration. * @return Returns 0 on success, or a non-zero error code on failure. */ int register_timeseries(std::string& device_id, @@ -145,4 +149,4 @@ class TsFileTreeWriter { } // namespace storage -#endif // WRITER_TSFILE_TREE_WRITER_H \ No newline at end of file +#endif // WRITER_TSFILE_TREE_WRITER_H diff --git a/cpp/src/writer/tsfile_writer.cc b/cpp/src/writer/tsfile_writer.cc index e72562136..a6cebfc09 100644 --- a/cpp/src/writer/tsfile_writer.cc +++ b/cpp/src/writer/tsfile_writer.cc @@ -305,7 +305,8 @@ int TsFileWriter::register_aligned_timeseries( MeasurementSchema* ms = new MeasurementSchema( measurement_schema.measurement_name_, measurement_schema.data_type_, measurement_schema.encoding_, measurement_schema.compression_type_); - int ret = register_timeseries(device_id, ms, true); + int ret = register_aligned_timeseries(device_id, + std::vector{ms}); if (ret != E_OK) { delete ms; } @@ -315,14 +316,37 @@ int TsFileWriter::register_aligned_timeseries( int TsFileWriter::register_aligned_timeseries( const std::string& device_id, const std::vector& measurement_schemas) { - int ret = E_OK; - for (auto it : measurement_schemas) { - ret = register_timeseries(device_id, it, true); - if (ret != E_OK) { + std::shared_ptr id = + std::make_shared(device_id); + // Fix the complete aligned schema at the first registration, like Java. + // schemas_ survives flush and is rebuilt from metadata during recovery. + if (schemas_.find(id) != schemas_.end() || measurement_schemas.empty()) { + return E_INVALID_ARG; + } + + std::unique_ptr group(new MeasurementSchemaGroup); + group->is_aligned_ = true; + for (auto* schema : measurement_schemas) { + if (schema == nullptr) { + return E_INVALID_ARG; + } + if (!group->measurement_schema_map_ + .insert(std::make_pair(schema->measurement_name_, schema)) + .second) { + return E_ALREADY_EXIST; + } + } + // Every registered column must participate from row 0. Publish the group + // and take ownership of its schemas only after the entire batch succeeds. + for (auto* schema : measurement_schemas) { + int ret = ensure_aligned_value_chunk_writer(schema); + if (RET_FAIL(ret)) { return ret; } } - return ret; + schemas_.insert(std::make_pair(id, group.get())); + group.release(); + return E_OK; } int TsFileWriter::register_timeseries( @@ -330,7 +354,7 @@ int TsFileWriter::register_timeseries( MeasurementSchema* ms = new MeasurementSchema( measurement_schema.measurement_name_, measurement_schema.data_type_, measurement_schema.encoding_, measurement_schema.compression_type_); - int ret = register_timeseries(device_id, ms, false); + int ret = register_timeseries(device_id, ms); if (ret != E_OK) { delete ms; } @@ -338,36 +362,23 @@ int TsFileWriter::register_timeseries( } int TsFileWriter::register_timeseries(const std::string& device_path, - MeasurementSchema* measurement_schema, - bool is_aligned) { + MeasurementSchema* measurement_schema) { + if (measurement_schema == nullptr) { + return E_INVALID_ARG; + } std::shared_ptr device_id = std::make_shared(device_path); DeviceSchemasMapIter device_iter = schemas_.find(device_id); if (device_iter != schemas_.end()) { MeasurementSchemaGroup* device_schema = device_iter->second; + // Non-aligned registration must not bypass the fixed aligned schema. + if (device_schema->is_aligned_) { + return E_INVALID_ARG; + } MeasurementSchemaMap& msm = device_schema->measurement_schema_map_; if (msm.find(measurement_schema->measurement_name_) != msm.end()) { return E_ALREADY_EXIST; } - if (device_schema->is_aligned_ && - device_schema->time_chunk_writer_ != nullptr && - device_schema->time_chunk_writer_->hasData()) { - // A column added now would have to be padded with the rows (and - // pages) that were already written, which the current page writer - // cannot express. Java does not allow an aligned device to be - // expanded at all; refuse loudly instead of silently writing a - // chunk group whose value column row counts diverge from the time - // column. - return E_INVALID_ARG; - } - // Aligned devices advance every registered measurement on every row, - // so the value chunk writer has to exist before the first row. - if (device_schema->is_aligned_) { - int ret = ensure_aligned_value_chunk_writer(measurement_schema); - if (RET_FAIL(ret)) { - return ret; - } - } MeasurementSchemaMapInsertResult ins_res = msm.insert(std::make_pair( measurement_schema->measurement_name_, measurement_schema)); if (UNLIKELY(!ins_res.second)) { @@ -375,14 +386,6 @@ int TsFileWriter::register_timeseries(const std::string& device_path, } } else { MeasurementSchemaGroup* ms_group = new MeasurementSchemaGroup; - ms_group->is_aligned_ = is_aligned; - if (is_aligned) { - int ret = ensure_aligned_value_chunk_writer(measurement_schema); - if (RET_FAIL(ret)) { - delete ms_group; - return ret; - } - } ms_group->measurement_schema_map_.insert(std::make_pair( measurement_schema->measurement_name_, measurement_schema)); schemas_.insert(std::make_pair(device_id, ms_group)); diff --git a/cpp/src/writer/tsfile_writer.h b/cpp/src/writer/tsfile_writer.h index 0737eb065..6ed68d8b2 100644 --- a/cpp/src/writer/tsfile_writer.h +++ b/cpp/src/writer/tsfile_writer.h @@ -69,9 +69,15 @@ class TsFileWriter { int register_timeseries( const std::string& device_path, const std::vector& measurement_schema_vec); + // Register the complete measurement list for an aligned device once. + // This overload fixes the device to a single measurement. Later aligned + // or non-aligned registration for that device returns E_INVALID_ARG, + // including before writing, after flush, and after recovery. int register_aligned_timeseries( const std::string& device_id, const MeasurementSchema& measurement_schema); + // The list must be nonempty with distinct, non-null schemas. Ownership + // transfers to the writer only when the entire registration succeeds. int register_aligned_timeseries( const std::string& device_id, const std::vector& measurement_schemas); @@ -186,8 +192,7 @@ class TsFileWriter { const Tablet& tablet, uint32_t start_idx = 0, uint32_t end_idx = UINT32_MAX); int register_timeseries(const std::string& device_path, - MeasurementSchema* measurement_schema, - bool is_aligned = false); + MeasurementSchema* measurement_schema); std::vector, int>> split_tablet_by_device(const Tablet& tablet); diff --git a/cpp/test/file/restorable_tsfile_io_writer_test.cc b/cpp/test/file/restorable_tsfile_io_writer_test.cc index 137f24f08..09296a996 100644 --- a/cpp/test/file/restorable_tsfile_io_writer_test.cc +++ b/cpp/test/file/restorable_tsfile_io_writer_test.cc @@ -485,6 +485,84 @@ TEST_F(RestorableTsFileIOWriterTest, AlignedTimeseriesRecoverAndWrite) { reader.close(); } +TEST_F(RestorableTsFileIOWriterTest, + RecoveredAlignedDeviceRejectsMeasurementRegistration) { + std::string device = "d1"; + { + TsFileWriter writer; + ASSERT_EQ(writer.open(file_name_, GetWriteCreateFlags(), 0666), E_OK); + std::vector schemas{ + new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED), + new MeasurementSchema("s1", INT64, PLAIN, UNCOMPRESSED)}; + ASSERT_EQ(writer.register_aligned_timeseries(device, schemas), E_OK); + TsRecord record(0, device); + record.add_point("s0", int64_t(100)); + // s1 is registered but entirely NULL before recovery. + ASSERT_EQ(writer.write_record_aligned(record), E_OK); + ASSERT_EQ(writer.flush(), E_OK); + ASSERT_EQ(writer.close(), E_OK); + } + CorruptCurrentFileTail(3); + + RestorableTsFileIOWriter recovery; + ASSERT_EQ(recovery.open(file_name_, true), E_OK); + ASSERT_TRUE(recovery.can_write()); + { + TsFileTreeWriter writer(&recovery); + for (const std::string& name : {"extra", "s0"}) { + auto* schema = + new MeasurementSchema(name, INT64, PLAIN, UNCOMPRESSED); + int ret = writer.register_timeseries( + device, std::vector{schema}); + EXPECT_EQ(ret, E_INVALID_ARG); + if (ret != E_OK) { + delete schema; + } + } + MeasurementSchema schema("nonaligned_extra", INT64, PLAIN, + UNCOMPRESSED); + EXPECT_EQ(writer.register_timeseries(device, &schema), E_INVALID_ARG); + TsRecord record(1, device); + record.add_point("s0", int64_t(101)); + record.add_point("s1", int64_t(201)); + ASSERT_EQ(writer.write(record), E_OK); + + // The restriction belongs to the recovered device, not the whole file. + std::string other_device = "d2"; + std::vector schemas{ + new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED)}; + ASSERT_EQ(writer.register_timeseries(other_device, schemas), E_OK); + TsRecord other(2, other_device); + other.add_point("s0", int64_t(302)); + ASSERT_EQ(writer.write(other), E_OK); + ASSERT_EQ(writer.flush(), E_OK); + ASSERT_EQ(writer.close(), E_OK); + } + + TsFileTreeReader reader; + ASSERT_EQ(reader.open(file_name_), E_OK); + ResultSet* result = nullptr; + ASSERT_EQ(reader.query({device}, {"s0", "s1"}, 0, 10, result), E_OK); + auto it = result->iterator(); + int row = 0; + while (it.hasNext()) { + RowRecord* record = it.next(); + EXPECT_EQ(record->get_timestamp(), row); + EXPECT_EQ(record->get_field(1)->type_, INT64); + EXPECT_EQ(record->get_field(1)->value_.lval_, 100 + row); + if (row == 0) { + EXPECT_EQ(record->get_field(2)->type_, NULL_TYPE); + } else { + EXPECT_EQ(record->get_field(2)->type_, INT64); + EXPECT_EQ(record->get_field(2)->value_.lval_, 201); + } + row++; + } + EXPECT_EQ(row, 2); + reader.destroy_query_data_set(result); + ASSERT_EQ(reader.close(), E_OK); +} + // ----------------------------------------------------------------------------- // Recovery + continued write with TsFileTableWriter (table model), then // read-back diff --git a/cpp/test/reader/tsfile_reader_test.cc b/cpp/test/reader/tsfile_reader_test.cc index 314d69c6d..e667b94da 100644 --- a/cpp/test/reader/tsfile_reader_test.cc +++ b/cpp/test/reader/tsfile_reader_test.cc @@ -253,16 +253,18 @@ TEST_P(MetadataReadLengthTest, RejectsIncompleteRangesBeforeParsing) { for (const bool aligned : {false, true}) { const std::string device = aligned ? "root.aligned" : "root.unaligned"; TsRecord record(100, device); + std::vector schemas; for (int column = 0; column < GetParam(); ++column) { const std::string name = "value" + std::to_string(column); - MeasurementSchema schema(name, INT32, PLAIN, UNCOMPRESSED); - ASSERT_EQ(aligned - ? tsfile_writer_->register_aligned_timeseries(device, - schema) - : tsfile_writer_->register_timeseries(device, schema), - E_OK); + schemas.push_back( + new MeasurementSchema(name, INT32, PLAIN, UNCOMPRESSED)); record.add_point(name, static_cast(42)); } + ASSERT_EQ( + aligned + ? tsfile_writer_->register_aligned_timeseries(device, schemas) + : tsfile_writer_->register_timeseries(device, schemas), + E_OK); ASSERT_EQ(aligned ? tsfile_writer_->write_record_aligned(record) : tsfile_writer_->write_record(record), E_OK); diff --git a/cpp/test/writer/tsfile_writer_test.cc b/cpp/test/writer/tsfile_writer_test.cc index b2d1e4cce..ea05f015b 100644 --- a/cpp/test/writer/tsfile_writer_test.cc +++ b/cpp/test/writer/tsfile_writer_test.cc @@ -480,18 +480,21 @@ TEST_F(TsFileWriterTest, WriteMultipleTabletsAlignedMultiFlush) { device_num, std::vector(measurement_num)); for (int i = 0; i < device_num; i++) { std::string device_name = "test_device" + std::to_string(i); + std::vector schemas; for (int j = 0; j < measurement_num; j++) { std::string measure_name = "measurement" + std::to_string(j); schema_vecs[i][j] = MeasurementSchema(measure_name, common::TSDataType::INT32, common::TSEncoding::PLAIN, common::CompressionType::UNCOMPRESSED); - tsfile_writer_->register_aligned_timeseries( - device_name, storage::MeasurementSchema( - measure_name, common::TSDataType::INT32, - common::TSEncoding::PLAIN, - common::CompressionType::UNCOMPRESSED)); + schemas.push_back( + new MeasurementSchema(measure_name, common::TSDataType::INT32, + common::TSEncoding::PLAIN, + common::CompressionType::UNCOMPRESSED)); } + ASSERT_EQ( + tsfile_writer_->register_aligned_timeseries(device_name, schemas), + E_OK); } for (int tablet_num = 0; tablet_num < max_tablet_num; tablet_num++) { @@ -1957,40 +1960,62 @@ TEST_F(TsFileWriterTest, AlignedRecordDuplicateMeasurementWritesOneRow) { ASSERT_EQ(reader.close(), E_OK); } -// New columns have to be padded with the rows (and pages) that were already -// written, which the aligned writer cannot express: reject the registration -// instead of writing a chunk group whose value column is shifted. -TEST_F(TsFileWriterTest, AlignedRegisterAfterWriteIsRejected) { - std::string device_name = "device_late_m"; - std::vector schemas; - schemas.push_back(new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED)); - tsfile_writer_->register_aligned_timeseries(device_name, schemas); +// Like Java, fix the complete aligned measurement list at registration, even +// before the first write. Flushing must not make the device extensible again. +class AlignedRegistrationTest + : public TsFileWriterTest, + public ::testing::WithParamInterface> {}; - TsRecord record(1, device_name); - record.add_point("s0", static_cast(1)); - ASSERT_EQ(tsfile_writer_->write_record_aligned(record), E_OK); +TEST_P(AlignedRegistrationTest, MeasurementsAreFixedAtRegistration) { + const bool single_measurement = std::get<0>(GetParam()); + const int stage = std::get<1>(GetParam()); + std::string device_name = "device_late_m"; + if (single_measurement) { + ASSERT_EQ(tsfile_writer_->register_aligned_timeseries( + device_name, + MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED)), + E_OK); + } else { + std::vector schemas{ + new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED), + new MeasurementSchema("s1", INT64, PLAIN, UNCOMPRESSED)}; + ASSERT_EQ( + tsfile_writer_->register_aligned_timeseries(device_name, schemas), + E_OK); + } + if (stage == 1) { + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + } + TsRecord first(0, device_name); + first.add_point("s0", int64_t(100)); + if (stage >= 2) { + ASSERT_EQ(tsfile_writer_->write_record_aligned(first), E_OK); + } + if (stage == 3) { + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + } // The writer only takes ownership of a schema when the registration // succeeds, so a rejected one has to be released by the caller. - MeasurementSchema* extra_schema = - new MeasurementSchema("s1", INT64, PLAIN, UNCOMPRESSED); - std::vector extra{extra_schema}; - EXPECT_EQ(tsfile_writer_->register_aligned_timeseries(device_name, extra), - E_INVALID_ARG); - delete extra_schema; - - // A second registration of the same measurement is still reported as a - // duplicate, not as a late registration. - MeasurementSchema* dup_schema = - new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED); - std::vector dup{dup_schema}; - EXPECT_EQ(tsfile_writer_->register_aligned_timeseries(device_name, dup), - E_ALREADY_EXIST); - delete dup_schema; + for (bool aligned : {true, false}) { + for (const std::string& name : {"extra", "s0"}) { + auto* schema = + new MeasurementSchema(name, INT64, PLAIN, UNCOMPRESSED); + std::vector extra{schema}; + int ret = aligned ? tsfile_writer_->register_aligned_timeseries( + device_name, extra) + : tsfile_writer_->register_timeseries(device_name, + extra); + EXPECT_EQ(ret, E_INVALID_ARG); + if (ret != E_OK) { + delete schema; + } + } + } // The schema-reference overloads own their internal copies even when the // registration is rejected. Exercise both wrappers under leak checking. - MeasurementSchema extra_value("s2", INT64, PLAIN, UNCOMPRESSED); + MeasurementSchema extra_value("extra_value", INT64, PLAIN, UNCOMPRESSED); EXPECT_EQ( tsfile_writer_->register_aligned_timeseries(device_name, extra_value), E_INVALID_ARG); @@ -1999,10 +2024,115 @@ TEST_F(TsFileWriterTest, AlignedRegisterAfterWriteIsRejected) { MeasurementSchema duplicate_value("s0", INT64, PLAIN, UNCOMPRESSED); EXPECT_EQ(tsfile_writer_->register_aligned_timeseries(device_name, duplicate_value), - E_ALREADY_EXIST); + E_INVALID_ARG); EXPECT_EQ(tsfile_writer_->register_timeseries(device_name, duplicate_value), - E_ALREADY_EXIST); + E_INVALID_ARG); + auto* groups = tsfile_writer_->get_schema_group_map(); + ASSERT_EQ(groups->size(), 1u); + ASSERT_EQ(groups->begin()->second->measurement_schema_map_.size(), + single_measurement ? 1u : 2u); + if (stage < 2) { + ASSERT_EQ(tsfile_writer_->write_record_aligned(first), E_OK); + } + TsRecord second(1, device_name); + second.add_point("s0", int64_t(101)); + if (!single_measurement) { + second.add_point("s1", int64_t(201)); + } + ASSERT_EQ(tsfile_writer_->write_record_aligned(second), E_OK); + + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); + + std::vector paths; + std::string s0("s0"), s1("s1"); + paths.emplace_back(device_name, s0); + if (!single_measurement) { + paths.emplace_back(device_name, s1); + } + TsFileReader reader; + ASSERT_EQ(reader.open(file_name_), E_OK); + ResultSet* result = nullptr; + ASSERT_EQ(reader.query(QueryExpression::create(paths, nullptr), result), + E_OK); + bool has_next = false; + int row = 0; + while (IS_SUCC(result->next(has_next)) && has_next) { + RowRecord* record = result->get_row_record(); + EXPECT_EQ(record->get_timestamp(), row); + EXPECT_EQ(field_to_string(record->get_field(1)), + std::to_string(100 + row)); + if (!single_measurement) { + EXPECT_EQ(field_to_string(record->get_field(2)), + row == 0 ? "NULL" : "201"); + } + row++; + } + EXPECT_EQ(row, 2); + reader.destroy_query_data_set(result); + ASSERT_EQ(reader.close(), E_OK); +} + +INSTANTIATE_TEST_SUITE_P(RegistrationStages, AlignedRegistrationTest, + ::testing::Combine(::testing::Bool(), + ::testing::Values(0, 1, 2, 3))); + +class AlignedRegistrationFailureTest + : public TsFileWriterTest, + public ::testing::WithParamInterface {}; + +TEST_P(AlignedRegistrationFailureTest, FailedBatchCanBeRetried) { + std::string device = "device_failed_registration"; + std::unique_ptr first( + new MeasurementSchema("s0", INT64, PLAIN, UNCOMPRESSED)); + std::unique_ptr second( + new MeasurementSchema("s1", INT64, PLAIN, UNCOMPRESSED)); + std::unique_ptr invalid(new MeasurementSchema( + GetParam() == 2 ? "s0" : "bad", INT64, DICTIONARY, UNCOMPRESSED)); + std::vector batch; + if (GetParam() == 1) { + batch = {first.get(), nullptr}; + } else if (GetParam() >= 2) { + batch = {first.get(), invalid.get()}; + } + // Empty list, null pointer, duplicate name, or unsupported encoding. + EXPECT_NE(tsfile_writer_->register_aligned_timeseries(device, batch), E_OK); + ASSERT_TRUE(tsfile_writer_->get_schema_group_map()->empty()); + + batch = {first.get(), second.get()}; + ASSERT_EQ(tsfile_writer_->register_aligned_timeseries(device, batch), E_OK); + first.release(); + second.release(); + TsRecord record(0, device); + record.add_point("s0", int64_t(100)); + record.add_point("s1", int64_t(200)); + ASSERT_EQ(tsfile_writer_->write_record_aligned(record), E_OK); + ASSERT_EQ(tsfile_writer_->flush(), E_OK); + ASSERT_EQ(tsfile_writer_->close(), E_OK); +} + +INSTANTIATE_TEST_SUITE_P(InvalidBatches, AlignedRegistrationFailureTest, + ::testing::Values(0, 1, 2, 3)); + +TEST_F(TsFileWriterTest, NonAlignedRegistrationRemainsIncremental) { + std::string device = "nonaligned_device"; + MeasurementSchema s0("s0", INT64, PLAIN, UNCOMPRESSED); + MeasurementSchema s1("s1", INT64, PLAIN, UNCOMPRESSED); + ASSERT_EQ(tsfile_writer_->register_timeseries(device, s0), E_OK); + ASSERT_EQ(tsfile_writer_->register_timeseries(device, s1), E_OK); + EXPECT_EQ(tsfile_writer_->register_timeseries(device, s0), E_ALREADY_EXIST); + EXPECT_EQ(tsfile_writer_->register_aligned_timeseries(device, s0), + E_INVALID_ARG); + EXPECT_EQ(tsfile_writer_->register_aligned_timeseries(device, s1), + E_INVALID_ARG); + auto* group = tsfile_writer_->get_schema_group_map()->begin()->second; + EXPECT_FALSE(group->is_aligned_); + EXPECT_EQ(group->measurement_schema_map_.size(), 2u); + TsRecord record(0, device); + record.add_point("s0", int64_t(100)); + record.add_point("s1", int64_t(200)); + ASSERT_EQ(tsfile_writer_->write_record(record), E_OK); ASSERT_EQ(tsfile_writer_->flush(), E_OK); ASSERT_EQ(tsfile_writer_->close(), E_OK); }