Skip to content
Open
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
117 changes: 46 additions & 71 deletions cpp/src/common/tsfile_common.cc
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
#include <algorithm>
#include <limits>
#include <map>
#include <utility>

#include "common/logger/elog.h"
#include "common/schema.h"
Expand Down Expand Up @@ -52,74 +53,56 @@ int TimeseriesIndex::add_chunk_meta(ChunkMeta* chunk_meta,

/* ================ TSMIterator ================ */
int TSMIterator::init() {
// sort chunk_group_meta_list_ : {[measurementA, offsetA1], [measurementB,
// offsetB1], [measurementA, offsetA2], [measurementB, offsetB2]} ->
// {[measurementA, offsetA1], [measurementA, offsetA2], [measurementB,
// offsetB1], [measurementB, offsetB2]}
tsm_chunk_meta_info_.clear();
const auto offset_less = [](const ChunkMeta* lhs, const ChunkMeta* rhs) {
return lhs->offset_of_chunk_header_ < rhs->offset_of_chunk_header_;
};
for (auto chunk_group_meta_iter = chunk_group_meta_list_.begin();
chunk_group_meta_iter != chunk_group_meta_list_.end();
chunk_group_meta_iter++) {
auto chunk_meta_list = chunk_group_meta_iter.get()->chunk_meta_list_;
// Use a map to group chunks by measurement_name_
std::map<common::String, std::vector<ChunkMeta*>> groups;
std::vector<common::String> order;
for (auto it = chunk_meta_list.begin(); it != chunk_meta_list.end();
it++) {
auto* chunk_meta = it.get();
if (groups.find(chunk_meta->measurement_name_) == groups.end()) {
order.push_back(chunk_meta->measurement_name_);
}
groups[chunk_meta->measurement_name_].push_back(chunk_meta);
}

// Sort each group of chunk metas by offset
for (auto it = groups.begin(); it != groups.end(); ++it) {
std::vector<ChunkMeta*>& group = it->second;
std::sort(group.begin(), group.end(),
[](ChunkMeta* a, ChunkMeta* b) {
return a->offset_of_chunk_header_ <
b->offset_of_chunk_header_;
});
ChunkGroupMeta* chunk_group_meta = chunk_group_meta_iter.get();
MeasurementChunkMetaMap chunk_group_map;
for (auto chunk_meta_iter = chunk_group_meta->chunk_meta_list_.begin();
chunk_meta_iter != chunk_group_meta->chunk_meta_list_.end();
chunk_meta_iter++) {
ChunkMeta* chunk_meta = chunk_meta_iter.get();
chunk_group_map[chunk_meta->measurement_name_].push_back(
chunk_meta);
}
// Clear and refill chunk_group_meta_list
chunk_group_meta_iter.get()->chunk_meta_list_.clear();
for (const auto& measurement_name : order) {
for (auto chunk_meta : groups[measurement_name]) {
chunk_group_meta_iter.get()->chunk_meta_list_.push_back(
chunk_meta);
for (auto& measurement_entry : chunk_group_map) {
auto& chunk_metas = measurement_entry.second;
if (!std::is_sorted(chunk_metas.begin(), chunk_metas.end(),
offset_less)) {
std::sort(chunk_metas.begin(), chunk_metas.end(), offset_less);
}
}
}

// FIXME empty list
chunk_group_meta_iter_ = chunk_group_meta_list_.begin();
while (chunk_group_meta_iter_ != chunk_group_meta_list_.end()) {
chunk_meta_iter_ =
chunk_group_meta_iter_.get()->chunk_meta_list_.begin();
std::map<common::String, std::vector<ChunkMeta*>> tmp;
while (chunk_meta_iter_ !=
chunk_group_meta_iter_.get()->chunk_meta_list_.end()) {
tmp[chunk_meta_iter_.get()->measurement_name_].emplace_back(
chunk_meta_iter_.get());
chunk_meta_iter_++;
}
if (!tmp.empty()) {
auto& merged =
tsm_chunk_meta_info_[chunk_group_meta_iter_.get()->device_id_];
for (auto& m_entry : tmp) {
auto& vec = merged[m_entry.first];
vec.insert(vec.end(), m_entry.second.begin(),
m_entry.second.end());
}
auto device_pos =
tsm_chunk_meta_info_.lower_bound(chunk_group_meta->device_id_);
IDeviceIDComparator comparator;
const bool device_exists =
device_pos != tsm_chunk_meta_info_.end() &&
!comparator(chunk_group_meta->device_id_, device_pos->first);
if (!device_exists) {
tsm_chunk_meta_info_.emplace_hint(device_pos,
chunk_group_meta->device_id_,
std::move(chunk_group_map));
continue;
}
Comment on lines +86 to 91

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Please skip empty chunk groups here and add a regression test for this case.


chunk_group_meta_iter_++;
}
if (!tsm_chunk_meta_info_.empty() &&
!tsm_chunk_meta_info_.begin()->second.empty()) {
tsm_measurement_iter_ = tsm_chunk_meta_info_.begin()->second.begin();
auto& measurement_map = device_pos->second;
for (auto& measurement_entry : chunk_group_map) {
auto& chunk_metas = measurement_entry.second;
auto& merged_chunk_metas = measurement_map[measurement_entry.first];
merged_chunk_metas.insert(merged_chunk_metas.end(),
chunk_metas.begin(), chunk_metas.end());
}
}

tsm_device_iter_ = tsm_chunk_meta_info_.begin();
if (tsm_device_iter_ != tsm_chunk_meta_info_.end()) {
tsm_measurement_iter_ = tsm_device_iter_->second.begin();
}
return E_OK;
}

Expand All @@ -131,8 +114,6 @@ int TSMIterator::get_next(std::shared_ptr<IDeviceID>& ret_device_id,
String& ret_measurement_name,
TimeseriesIndex& ret_ts_index) {
int ret = E_OK;
SimpleList<ChunkMeta*> chunk_meta_list_of_this_ts(
1024, MOD_TIMESERIES_INDEX_OBJ); // FIXME
if (tsm_measurement_iter_ == tsm_device_iter_->second.end()) {
tsm_device_iter_++;
if (!has_next()) {
Expand All @@ -143,15 +124,13 @@ int TSMIterator::get_next(std::shared_ptr<IDeviceID>& ret_device_id,
}
ret_device_id = tsm_device_iter_->first;
ret_measurement_name.shallow_copy_from(tsm_measurement_iter_->first);
for (auto meta : tsm_measurement_iter_->second) {
chunk_meta_list_of_this_ts.push_back(meta);
}
if (chunk_meta_list_of_this_ts.size() == 0) {
const std::vector<ChunkMeta*>& chunk_metas = tsm_measurement_iter_->second;
if (chunk_metas.empty()) {
return E_TSFILE_WRITER_META_ERR;
}

const bool multi_chunks = chunk_meta_list_of_this_ts.size() > 1;
ChunkMeta* first_chunk_meta = chunk_meta_list_of_this_ts.front();
const bool multi_chunks = chunk_metas.size() > 1;
ChunkMeta* first_chunk_meta = chunk_metas.front();
const char meta_type = (multi_chunks ? 1 : 0) | (first_chunk_meta->mask_);
const TSDataType data_type = first_chunk_meta->data_type_;

Expand All @@ -160,13 +139,9 @@ int TSMIterator::get_next(std::shared_ptr<IDeviceID>& ret_device_id,
ret_ts_index.set_data_type(data_type);
ret_ts_index.init_statistic(data_type);

SimpleList<ChunkMeta*>::Iterator ts_chunk_meta_iter =
chunk_meta_list_of_this_ts.begin();
for (;
IS_SUCC(ret) && ts_chunk_meta_iter != chunk_meta_list_of_this_ts.end();
ts_chunk_meta_iter++) {
ChunkMeta* chunk_meta = ts_chunk_meta_iter.get();
for (ChunkMeta* chunk_meta : chunk_metas) {
if (RET_FAIL(ret_ts_index.add_chunk_meta(chunk_meta, multi_chunks))) {
break;
}
}
if (IS_SUCC(ret)) {
Expand Down
24 changes: 10 additions & 14 deletions cpp/src/common/tsfile_common.h
Original file line number Diff line number Diff line change
Expand Up @@ -695,9 +695,7 @@ class TSMIterator {
public:
explicit TSMIterator(
common::SimpleList<ChunkGroupMeta*>& chunk_group_meta_list)
: chunk_group_meta_list_(chunk_group_meta_list),
chunk_group_meta_iter_(),
chunk_meta_iter_() {}
: chunk_group_meta_list_(chunk_group_meta_list) {}

// sort => iterate
int init();
Expand All @@ -707,24 +705,22 @@ class TSMIterator {
TimeseriesIndex& ret_ts_index);

private:
using MeasurementChunkMetaMap =
std::map<common::String, std::vector<ChunkMeta*>>;
using DeviceChunkMetaMap =
std::map<std::shared_ptr<IDeviceID>, MeasurementChunkMetaMap,
IDeviceIDComparator>;

common::SimpleList<ChunkGroupMeta*>& chunk_group_meta_list_;
common::SimpleList<ChunkGroupMeta*>::Iterator chunk_group_meta_iter_;
common::SimpleList<ChunkMeta*>::Iterator chunk_meta_iter_;

// timeseries measurenemnt chunk meta info
std::map<std::shared_ptr<IDeviceID>,
std::map<common::String, std::vector<ChunkMeta*>>,
IDeviceIDComparator>
tsm_chunk_meta_info_;
DeviceChunkMetaMap tsm_chunk_meta_info_;

// device iterator
std::map<std::shared_ptr<IDeviceID>,
std::map<common::String, std::vector<ChunkMeta*>>,
IDeviceIDComparator>::iterator tsm_device_iter_;
DeviceChunkMetaMap::iterator tsm_device_iter_;

// measurement iterator
std::map<common::String, std::vector<ChunkMeta*>>::iterator
tsm_measurement_iter_;
MeasurementChunkMetaMap::iterator tsm_measurement_iter_;
};

/* =============== TsFile Index ================ */
Expand Down
40 changes: 16 additions & 24 deletions cpp/src/file/tsfile_io_writer.cc
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,8 @@ using namespace common;

namespace storage {

const uint32_t TsFileIOWriter::WRITE_STREAM_PAGE_SIZE;

#if 0
#define OFFSET_DEBUG(msg) \
std::cout << "OFFSET_DEBUG: " << msg \
Expand Down Expand Up @@ -379,7 +381,7 @@ int TsFileIOWriter::write_file_index() {
std::shared_ptr<IMetaIndexEntry> meta_index_entry = nullptr;
std::shared_ptr<MetaIndexNode> cur_index_node = nullptr;
SimpleList<std::shared_ptr<MetaIndexNode>>* cur_index_node_queue = nullptr;
DeviceNodeMap device_map;
std::map<std::string, DeviceNodeMap> table_device_nodes_map;

TSMIterator tsm_iter(chunk_group_meta_list_);

Expand Down Expand Up @@ -410,9 +412,11 @@ int TsFileIOWriter::write_file_index() {
if (prev_device_id != nullptr) {
if (RET_FAIL(add_cur_index_node_to_queue(
cur_index_node, cur_index_node_queue))) {
} else if (RET_FAIL(add_device_node(device_map, prev_device_id,
cur_index_node_queue,
writing_mm))) {
} else if (RET_FAIL(add_device_node(
table_device_nodes_map[prev_device_id
->get_table_name()],
prev_device_id, cur_index_node_queue,
writing_mm))) {
}
}
if (IS_SUCC(ret)) {
Expand Down Expand Up @@ -477,26 +481,16 @@ int TsFileIOWriter::write_file_index() {
ASSERT(cur_index_node_queue != nullptr);
if (RET_FAIL(add_cur_index_node_to_queue(cur_index_node,
cur_index_node_queue))) {
} else if (RET_FAIL(add_device_node(device_map, prev_device_id,
cur_index_node_queue,
writing_mm))) {
} else if (RET_FAIL(add_device_node(
table_device_nodes_map[prev_device_id->get_table_name()],
prev_device_id, cur_index_node_queue, writing_mm))) {
}
}

if (IS_SUCC(ret)) {
TsFileMeta tsfile_meta;
tsfile_meta.meta_offset_ = meta_offset;
tsfile_meta.bloom_filter_ = &filter;
// split device by table
std::map<std::string, DeviceNodeMap> table_device_nodes_map;
for (const auto& entry : device_map) {
std::string table_name = entry.first->get_table_name();
auto& table_map = table_device_nodes_map[table_name];
if (table_map.empty() ||
table_map.find(entry.first) == table_map.end()) {
table_map[entry.first] = entry.second;
}
}
std::map<std::string, std::shared_ptr<MetaIndexNode>> table_nodes_map;
for (auto& entry : table_device_nodes_map) {
auto meta_index_node =
Expand Down Expand Up @@ -730,20 +724,18 @@ int TsFileIOWriter::add_device_node(
FileIndexWritingMemManager& wmm) {
ASSERT(measurement_index_node_queue->size() > 0);
int ret = E_OK;
auto find_iter = device_map.find(device_id);
if (find_iter != device_map.end()) {
auto insert_pos = device_map.lower_bound(device_id);
IDeviceIDComparator comparator;
if (insert_pos != device_map.end() &&
!comparator(device_id, insert_pos->first)) {
return E_ALREADY_EXIST;
}

std::shared_ptr<MetaIndexNode> root = nullptr;
if (RET_FAIL(generate_root(measurement_index_node_queue, root,
INTERNAL_MEASUREMENT, wmm))) {
} else {
std::pair<DeviceNodeMapIterator, bool> ins_res =
device_map.insert(std::make_pair(device_id, root));
if (!ins_res.second) {
ASSERT(false);
}
device_map.emplace_hint(insert_pos, device_id, root);
}
return ret;
}
Expand Down
3 changes: 2 additions & 1 deletion cpp/src/file/tsfile_io_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,8 @@ class TsFileIOWriter {
typedef DeviceNodeMap::iterator DeviceNodeMapIterator;

public:
static const uint32_t WRITE_STREAM_PAGE_SIZE = 512; // FIXME
static const uint32_t WRITE_STREAM_PAGE_SIZE = 16 * 1024;

public:
TsFileIOWriter()
: meta_allocator_(),
Expand Down
Loading
Loading