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 41a485c00..a6cebfc09 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" @@ -303,20 +305,48 @@ 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_aligned_timeseries(device_id, + std::vector{ms}); + if (ret != E_OK) { + delete ms; + } + return ret; } 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( @@ -324,18 +354,31 @@ 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); + if (ret != E_OK) { + delete ms; + } + return ret; } 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()) { - MeasurementSchemaMap& msm = - device_iter->second->measurement_schema_map_; + 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; + } MeasurementSchemaMapInsertResult ins_res = msm.insert(std::make_pair( measurement_schema->measurement_name_, measurement_schema)); if (UNLIKELY(!ins_res.second)) { @@ -343,7 +386,6 @@ int TsFileWriter::register_timeseries(const std::string& device_path, } } else { MeasurementSchemaGroup* ms_group = new MeasurementSchemaGroup; - ms_group->is_aligned_ = is_aligned; 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 +558,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 +610,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 +853,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 +909,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 +925,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; @@ -881,15 +974,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]; @@ -984,15 +1082,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 +1137,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 +1171,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..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); @@ -128,6 +134,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, @@ -178,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/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..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 @@ -1061,3 +1139,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/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 3b9dae92a..ea05f015b 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 @@ -479,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++) { @@ -1737,3 +1741,637 @@ 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); +} + +// 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> {}; + +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. + 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("extra_value", 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_INVALID_ARG); + EXPECT_EQ(tsfile_writer_->register_timeseries(device_name, duplicate_value), + 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); +} + +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. +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); +}