Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions cpp/src/writer/tsfile_tree_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -145,4 +149,4 @@ class TsFileTreeWriter {

} // namespace storage

#endif // WRITER_TSFILE_TREE_WRITER_H
#endif // WRITER_TSFILE_TREE_WRITER_H
257 changes: 196 additions & 61 deletions cpp/src/writer/tsfile_writer.cc

Large diffs are not rendered by default.

17 changes: 15 additions & 2 deletions cpp/src/writer/tsfile_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -69,9 +69,15 @@ class TsFileWriter {
int register_timeseries(
const std::string& device_path,
const std::vector<MeasurementSchema*>& 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<MeasurementSchema*>& measurement_schemas);
Expand Down Expand Up @@ -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<ValueChunkWriter*>& value_chunk_writers,
Expand Down Expand Up @@ -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<std::pair<std::shared_ptr<IDeviceID>, int>>
split_tablet_by_device(const Tablet& tablet);

Expand Down
35 changes: 35 additions & 0 deletions cpp/src/writer/value_chunk_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -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_; }
Expand Down
20 changes: 20 additions & 0 deletions cpp/src/writer/value_page_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -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() {
Expand Down
204 changes: 204 additions & 0 deletions cpp/test/file/restorable_tsfile_io_writer_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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<MeasurementSchema*> 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<MeasurementSchema*>{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<MeasurementSchema*> 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
Expand Down Expand Up @@ -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<string> even_names = {"s1", "s2", "s3"};
vector<string> odd_names = {"s4", "s5", "s6"};
{
TsFileWriter tw;
ASSERT_EQ(tw.open(file_name_, GetWriteCreateFlags(), 0666), E_OK);
std::vector<MeasurementSchema*> 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<int32_t>(i));
record.add_point(even_names[2], "even");
} else {
record.add_point(odd_names[0], static_cast<int64_t>(i));
record.add_point(odd_names[1], static_cast<float>(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<int32_t>(i));
record.add_point(even_names[2], "even");
} else {
record.add_point(odd_names[0], static_cast<int64_t>(i));
record.add_point(odd_names[1], static_cast<float>(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<std::string, storage::ITimeseriesIndex*> 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<string> measurement_names = {"s1", "s2", "s3", "s4", "s5", "s6"};
ASSERT_EQ(CountTreeReaderRows(reader, measurement_names), 20);

ResultSet* result_set = nullptr;
vector<string> 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();
}
14 changes: 8 additions & 6 deletions cpp/test/reader/tsfile_reader_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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<MeasurementSchema*> 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<int32_t>(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);
Expand Down
Loading
Loading