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
12 changes: 2 additions & 10 deletions cpp/src/cwrapper/tsfile_cwrapper.cc
Original file line number Diff line number Diff line change
Expand Up @@ -318,11 +318,7 @@ ERRNO tsfile_writer_close(TsFileWriter writer) {
return common::E_OK;
}
auto* w = static_cast<storage::TsFileTableWriter*>(writer);
int ret = w->flush();
if (ret != common::E_OK) {
return ret;
}
ret = w->close();
int ret = w->close();
if (ret != common::E_OK) {
return ret;
}
Expand Down Expand Up @@ -2300,11 +2296,7 @@ ERRNO _tsfile_writer_write_ts_record(TsFileWriter writer, TsRecord data) {

ERRNO _tsfile_writer_close(TsFileWriter writer) {
auto* w = static_cast<storage::TsFileWriter*>(writer);
int ret = w->flush();
if (ret != common::E_OK) {
return ret;
}
ret = w->close();
int ret = w->close();
if (ret != common::E_OK) {
return ret;
}
Expand Down
3 changes: 2 additions & 1 deletion cpp/src/cwrapper/tsfile_cwrapper.h
Original file line number Diff line number Diff line change
Expand Up @@ -496,7 +496,8 @@ TsFileWriter tsfile_writer_new_with_memory_threshold(WriteFile file,
TsFileReader tsfile_reader_new(const char* pathname, ERRNO* err_code);

/**
* @brief Releases resources associated with a TsFileWriter.
* @brief Flushes pending data, finalizes the TsFile, and releases resources
* associated with a TsFileWriter.
Comment on lines +499 to +500
*
* @param writer [in] Writer handle obtained from tsfile_writer_new().
* After call: handle becomes invalid and must not be reused.
Expand Down
2 changes: 1 addition & 1 deletion cpp/src/writer/tsfile_table_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -125,7 +125,7 @@ class TsFileTableWriter {
int add_tsfile_property(const std::string& key,
const std::vector<uint8_t>& value);
/**
* Closes the writer and releases any resources held by it.
* Flushes pending data, finalizes the file, and releases writer resources.
* After calling this method, no further operations should be performed on
* this instance.
*
Expand Down
4 changes: 4 additions & 0 deletions cpp/src/writer/tsfile_writer.cc
Original file line number Diff line number Diff line change
Expand Up @@ -1986,6 +1986,10 @@ int TsFileWriter::flush_chunk_group(MeasurementSchemaGroup* chunk_group,

int TsFileWriter::close() {
if (UNLIKELY(unrecoverable_)) return E_DATA_INCONSISTENCY;
int ret = flush();
if (ret != E_OK) {
return ret;
}
return io_writer_->end_file();
Comment on lines +1989 to 1993
}

Expand Down
5 changes: 3 additions & 2 deletions cpp/src/writer/tsfile_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -115,8 +115,9 @@ class TsFileWriter {
int flush();

/*
* Flush file index part of the whole file (it may be flushed many times
* before close, the index part should cover all data in disk file).
* Flushes remaining buffered data, writes the file index and footer, and
* closes the file. Any flush failure is returned without finalizing the
* file, so the caller can handle the error and retry if possible.
*/
int close();

Expand Down
22 changes: 21 additions & 1 deletion cpp/test/writer/table_view/tsfile_writer_table_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -139,8 +139,28 @@ TEST_F(TsFileWriterTableTest, WriteTableTest) {
std::make_shared<TsFileTableWriter>(&write_file_, table_schema);
auto tablet = gen_tablet(table_schema, 0, 1);
ASSERT_EQ(tsfile_table_writer_->write_table(tablet), common::E_OK);
ASSERT_EQ(tsfile_table_writer_->flush(), common::E_OK);
ASSERT_EQ(tsfile_table_writer_->close(), common::E_OK);

TsFileReader reader;
ASSERT_EQ(reader.open(file_name_), common::E_OK);
ResultSet* result_set = nullptr;
ASSERT_EQ(reader.query(table_schema->get_table_name(), {"s0"}, 0, INT32_MAX,
result_set),
common::E_OK);
auto* table_result_set = static_cast<TableResultSet*>(result_set);
bool has_next = false;
int64_t row_count = 0;
while (true) {
ASSERT_EQ(table_result_set->next(has_next), common::E_OK);
if (!has_next) {
break;
}
++row_count;
}
EXPECT_EQ(row_count, 10);
table_result_set->close();
reader.destroy_query_data_set(table_result_set);
ASSERT_EQ(reader.close(), common::E_OK);
delete table_schema;
}

Expand Down
43 changes: 43 additions & 0 deletions cpp/test/writer/tsfile_writer_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -871,6 +871,49 @@ TEST_F(TsFileWriterTest, FlushWithoutWriteAfterRegisterTS) {
ASSERT_EQ(tsfile_writer_->close(), E_OK);
}

TEST_F(TsFileWriterTest, CloseFlushesDataWrittenAfterLastFlush) {
const std::string device_path = "device_close_flush";
const std::string measurement_name = "value";
ASSERT_EQ(tsfile_writer_->register_timeseries(
device_path, storage::MeasurementSchema(
measurement_name, common::TSDataType::INT64,
common::TSEncoding::PLAIN,
common::CompressionType::UNCOMPRESSED)),
E_OK);

TsRecord first_record(100, device_path);
first_record.add_point(measurement_name, static_cast<int64_t>(1));
ASSERT_EQ(tsfile_writer_->write_record(first_record), E_OK);
ASSERT_EQ(tsfile_writer_->flush(), E_OK);

TsRecord final_record(101, device_path);
final_record.add_point(measurement_name, static_cast<int64_t>(2));
ASSERT_EQ(tsfile_writer_->write_record(final_record), E_OK);
ASSERT_EQ(tsfile_writer_->close(), E_OK);

TsFileReader reader;
ASSERT_EQ(reader.open(file_name_), E_OK);
ResultSet* result_set = nullptr;
std::vector<std::string> select_list = {device_path + "." +
measurement_name};
ASSERT_EQ(reader.query(select_list, 100, 102, result_set), E_OK);
auto* query_result = static_cast<QDSWithoutTimeGenerator*>(result_set);
bool has_next = false;
ASSERT_EQ(query_result->next(has_next), E_OK);
ASSERT_TRUE(has_next);
EXPECT_EQ(query_result->get_value<int64_t>(1), 100);
EXPECT_EQ(query_result->get_value<int64_t>(2), 1);
ASSERT_EQ(query_result->next(has_next), E_OK);
ASSERT_TRUE(has_next);
EXPECT_EQ(query_result->get_value<int64_t>(1), 101);
EXPECT_EQ(query_result->get_value<int64_t>(2), 2);
ASSERT_EQ(query_result->next(has_next), E_OK);
EXPECT_FALSE(has_next);

reader.destroy_query_data_set(result_set);
ASSERT_EQ(reader.close(), E_OK);
}

TEST_F(TsFileWriterTest, WriteAlignedTimeseries) {
int measurement_num = 100, row_num = 150;
std::string device_name = "device";
Expand Down
Loading