Skip to content

Commit cb1c9b7

Browse files
committed
Make FragmentInfo::getNumTuples a getter
Move the logic for synthetic fragment number of tuples into FragmentInfo, this will enable computing it only when necessary.
1 parent d2e5a1d commit cb1c9b7

5 files changed

Lines changed: 57 additions & 30 deletions

File tree

Fragmenter/Fragmenter.h

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -55,7 +55,13 @@ struct InsertData {
5555

5656
class FragmentInfo {
5757
public:
58-
FragmentInfo() : fragmentId(-1), numTuples(0), shadowNumTuples(0), resultSet(nullptr) {}
58+
FragmentInfo()
59+
: fragmentId(-1),
60+
shadowNumTuples(0),
61+
resultSet(nullptr),
62+
numTuples(0),
63+
synthesizedNumTuplesIsValid(false),
64+
synthesizedMetadataIsValid(false) {}
5965

6066
void setChunkMetadataMap(const std::map<int, ChunkMetadata>& chunkMetadataMap) {
6167
this->chunkMetadataMap = chunkMetadataMap;
@@ -67,15 +73,24 @@ class FragmentInfo {
6773

6874
const std::map<int, ChunkMetadata>& getChunkMetadataMapPhysical() const { return chunkMetadataMap; }
6975

76+
size_t getNumTuples() const;
77+
78+
size_t getPhysicalNumTuples() const { return numTuples; }
79+
80+
void setPhysicalNumTuples(const size_t physNumTuples) { numTuples = physNumTuples; }
81+
7082
int fragmentId;
71-
size_t numTuples;
7283
size_t shadowNumTuples;
7384
std::vector<int> deviceIds;
7485
std::map<int, ChunkMetadata> shadowChunkMetadataMap;
7586
mutable ResultRows* resultSet;
87+
mutable std::shared_ptr<std::mutex> resultSetMutex;
7688

7789
private:
90+
mutable size_t numTuples;
7891
mutable std::map<int, ChunkMetadata> chunkMetadataMap;
92+
mutable bool synthesizedNumTuplesIsValid;
93+
mutable bool synthesizedMetadataIsValid;
7994
};
8095

8196
/**

Fragmenter/InsertOrderFragmenter.cpp

Lines changed: 10 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -76,14 +76,14 @@ void InsertOrderFragmenter::getChunkMetadata() {
7676
maxFragmentId_ = curFragmentId;
7777
fragmentInfoVec_.push_back(FragmentInfo());
7878
fragmentInfoVec_.back().fragmentId = curFragmentId;
79-
fragmentInfoVec_.back().numTuples = chunkIt->second.numElements;
80-
numTuples_ += fragmentInfoVec_.back().numTuples;
79+
fragmentInfoVec_.back().setPhysicalNumTuples(chunkIt->second.numElements);
80+
numTuples_ += fragmentInfoVec_.back().getPhysicalNumTuples();
8181
for (const auto levelSize : dataMgr_->levelSizes_) {
8282
fragmentInfoVec_.back().deviceIds.push_back(curFragmentId % levelSize);
8383
}
84-
fragmentInfoVec_.back().shadowNumTuples = fragmentInfoVec_.back().numTuples;
84+
fragmentInfoVec_.back().shadowNumTuples = fragmentInfoVec_.back().getPhysicalNumTuples();
8585
} else {
86-
if (chunkIt->second.numElements != fragmentInfoVec_.back().numTuples) {
86+
if (chunkIt->second.numElements != fragmentInfoVec_.back().getPhysicalNumTuples()) {
8787
throw std::runtime_error("Inconsistency in num tuples within fragment");
8888
}
8989
}
@@ -134,7 +134,7 @@ void InsertOrderFragmenter::dropFragmentsToSize(const size_t maxRows) {
134134
size_t targetRows = maxRows * DROP_FRAGMENT_FACTOR;
135135
while (numTuples_ > targetRows) {
136136
assert(fragmentInfoVec_.size() > 0);
137-
size_t numFragTuples = fragmentInfoVec_[0].numTuples;
137+
size_t numFragTuples = fragmentInfoVec_[0].getPhysicalNumTuples();
138138
dropFragIds.push_back(fragmentInfoVec_[0].fragmentId);
139139
fragmentInfoVec_.pop_front();
140140
assert(numTuples_ >= numFragTuples);
@@ -256,13 +256,13 @@ void InsertOrderFragmenter::insertData(const InsertData& insertDataStruct) {
256256
delete[] rowIdData;
257257
}
258258

259-
currentFragment->shadowNumTuples = fragmentInfoVec_.back().numTuples + numRowsToInsert;
259+
currentFragment->shadowNumTuples = fragmentInfoVec_.back().getPhysicalNumTuples() + numRowsToInsert;
260260
numRowsLeft -= numRowsToInsert;
261261
numRowsInserted += numRowsToInsert;
262262
}
263263
mapd_unique_lock<mapd_shared_mutex> writeLock(fragmentInfoMutex_);
264264
for (auto partIt = fragmentInfoVec_.begin() + startFragment; partIt != fragmentInfoVec_.end(); ++partIt) {
265-
partIt->numTuples = partIt->shadowNumTuples;
265+
partIt->setPhysicalNumTuples(partIt->shadowNumTuples);
266266
partIt->setChunkMetadataMap(partIt->shadowChunkMetadataMap);
267267
}
268268
numTuples_ += insertDataStruct.numRows;
@@ -277,7 +277,7 @@ FragmentInfo* InsertOrderFragmenter::createNewFragment(const Data_Namespace::Mem
277277
FragmentInfo newFragmentInfo;
278278
newFragmentInfo.fragmentId = maxFragmentId_;
279279
newFragmentInfo.shadowNumTuples = 0;
280-
newFragmentInfo.numTuples = 0;
280+
newFragmentInfo.setPhysicalNumTuples(0);
281281
for (const auto levelSize : dataMgr_->levelSizes_) {
282282
newFragmentInfo.deviceIds.push_back(newFragmentInfo.fragmentId % levelSize);
283283
}
@@ -307,15 +307,15 @@ TableInfo InsertOrderFragmenter::getFragmentsForQuery() {
307307
queryInfo.numTuples = 0;
308308
auto partIt = queryInfo.fragments.begin();
309309
while (partIt != queryInfo.fragments.end()) {
310-
if (partIt->numTuples == 0) {
310+
if (partIt->getPhysicalNumTuples() == 0) {
311311
// this means that a concurrent insert query inserted tuples into a new fragment but when the query came in we
312312
// didn't have this fragment.
313313
// To make sure we don't mess up the executor we delete this
314314
// fragment from the metadatamap (fixes earlier bug found
315315
// 2015-05-08)
316316
partIt = queryInfo.fragments.erase(partIt);
317317
} else {
318-
queryInfo.numTuples += partIt->numTuples;
318+
queryInfo.numTuples += partIt->getPhysicalNumTuples();
319319
++partIt;
320320
}
321321
}

QueryEngine/Execute.cpp

Lines changed: 12 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -4011,11 +4011,11 @@ size_t compute_buffer_entry_guess(const std::vector<InputTableInfo>& query_infos
40114011
size_t max_groups_buffer_entry_guess = 1;
40124012
for (const auto& query_info : query_infos) {
40134013
CHECK(!query_info.info.fragments.empty());
4014-
auto it =
4015-
std::max_element(query_info.info.fragments.begin(),
4016-
query_info.info.fragments.end(),
4017-
[](const FragmentInfo& f1, const FragmentInfo& f2) { return f1.numTuples < f2.numTuples; });
4018-
max_groups_buffer_entry_guess *= it->numTuples;
4014+
auto it = std::max_element(
4015+
query_info.info.fragments.begin(),
4016+
query_info.info.fragments.end(),
4017+
[](const FragmentInfo& f1, const FragmentInfo& f2) { return f1.getNumTuples() < f2.getNumTuples(); });
4018+
max_groups_buffer_entry_guess *= it->getNumTuples();
40194019
}
40204020
return max_groups_buffer_entry_guess;
40214021
}
@@ -4332,7 +4332,7 @@ Executor::ExecutionDispatch::ExecutionDispatch(Executor* executor,
43324332
render_allocator_map_(render_allocator_map) {
43334333
all_frag_row_offsets_.resize(query_infos.front().info.fragments.size() + 1);
43344334
for (size_t i = 1; i <= query_infos.front().info.fragments.size(); ++i) {
4335-
all_frag_row_offsets_[i] = all_frag_row_offsets_[i - 1] + query_infos_.front().info.fragments[i - 1].numTuples;
4335+
all_frag_row_offsets_[i] = all_frag_row_offsets_[i - 1] + query_infos_.front().info.fragments[i - 1].getNumTuples();
43364336
}
43374337
all_fragment_results_.reserve(query_infos_.front().info.fragments.size());
43384338
}
@@ -4686,7 +4686,7 @@ const int8_t* Executor::ExecutionDispatch::getAllScanColumnFrags(
46864686
Data_Namespace::CPU_LEVEL,
46874687
int(0));
46884688
column_frags.push_back(boost::make_unique<ColumnarResults>(
4689-
row_set_mem_owner_, col_buffer, fragment.numTuples, chunk_meta_it->second.sqlType));
4689+
row_set_mem_owner_, col_buffer, fragment.getNumTuples(), chunk_meta_it->second.sqlType));
46904690
}
46914691
column_it->second = ColumnarResults::mergeResults(row_set_mem_owner_, column_frags);
46924692
}
@@ -4729,7 +4729,7 @@ std::vector<const ColumnarResults*> Executor::ExecutionDispatch::getAllScanColum
47294729
frags_it->second.insert(
47304730
std::make_pair(CacheKey{frag_id},
47314731
boost::make_unique<ColumnarResults>(
4732-
row_set_mem_owner_, col_buffer, fragment.numTuples, chunk_meta_it->second.sqlType)));
4732+
row_set_mem_owner_, col_buffer, fragment.getNumTuples(), chunk_meta_it->second.sqlType)));
47334733
}
47344734
}
47354735
CHECK(frags_it != columnarized_ref_table_cache_.end());
@@ -4855,7 +4855,7 @@ const int8_t* Executor::ExecutionDispatch::getColumn(
48554855
Data_Namespace::CPU_LEVEL,
48564856
device_id);
48574857
ColumnarResults ref_values(
4858-
row_set_mem_owner_, col_buffer, fragment.numTuples, chunk_meta_it->second.sqlType);
4858+
row_set_mem_owner_, col_buffer, fragment.getNumTuples(), chunk_meta_it->second.sqlType);
48594859
frag_id_to_result.insert(
48604860
std::make_pair(sub_key,
48614861
ColumnarResults::createIndexedResults(
@@ -5338,7 +5338,7 @@ void Executor::dispatchFragments(
53385338
frag_list_idx % context_count,
53395339
rowid_lookup_key));
53405340
++frag_list_idx;
5341-
if (is_sample_query(ra_exe_unit) && fragment.numTuples >= ra_exe_unit.scan_limit) {
5341+
if (is_sample_query(ra_exe_unit) && fragment.getNumTuples() >= ra_exe_unit.scan_limit) {
53425342
break;
53435343
}
53445344
}
@@ -5397,7 +5397,7 @@ std::map<size_t, std::vector<uint64_t>> Executor::getAllFragOffsets(
53975397
std::vector<uint64_t> frag_offsets(fragments.size(), 0);
53985398
for (size_t i = 0, off = 0; i < fragments.size(); ++i) {
53995399
frag_offsets[i] = off;
5400-
off += fragments[i].numTuples;
5400+
off += fragments[i].getNumTuples();
54015401
}
54025402
tab_id_to_frag_offsets.insert(std::make_pair(desc.getTableId(), frag_offsets));
54035403
}
@@ -5453,7 +5453,7 @@ Executor::FetchResult Executor::fetchChunks(const ExecutionDispatch& execution_d
54535453
CHECK(fragments_it != all_tables_fragments.end());
54545454
const auto& fragments = *fragments_it->second;
54555455
const auto& fragment = fragments[frag_id];
5456-
num_rows.push_back(fragment.numTuples);
5456+
num_rows.push_back(fragment.getNumTuples());
54575457
const auto frag_offsets_it = tab_id_to_frag_offsets.find(input_descs[tab_idx].getTableId());
54585458
CHECK(frag_offsets_it != tab_id_to_frag_offsets.end());
54595459
const auto& offsets = frag_offsets_it->second;

QueryEngine/InputMetadata.cpp

Lines changed: 17 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -153,9 +153,9 @@ Fragmenter_Namespace::TableInfo synthesize_table_info(const RowSetPtr& rows) {
153153
result.resize(1);
154154
auto& fragment = result.front();
155155
fragment.fragmentId = 0;
156-
fragment.numTuples = row_count;
157156
fragment.deviceIds.resize(3);
158157
fragment.resultSet = rows.get();
158+
fragment.resultSetMutex.reset(new std::mutex());
159159
}
160160
Fragmenter_Namespace::TableInfo table_info;
161161
table_info.fragments = result;
@@ -171,9 +171,9 @@ Fragmenter_Namespace::TableInfo synthesize_table_info(const IterTabPtr& table) {
171171
for (size_t i = 0; i < table->fragCount(); ++i) {
172172
auto& fragment = table_info.fragments[i];
173173
fragment.fragmentId = i;
174-
fragment.numTuples = table->getFragAt(i).row_count;
174+
fragment.setPhysicalNumTuples(table->getFragAt(i).row_count);
175175
fragment.deviceIds.resize(3);
176-
total_row_count += fragment.numTuples;
176+
total_row_count += fragment.getPhysicalNumTuples();
177177
}
178178
}
179179

@@ -247,9 +247,21 @@ std::vector<InputTableInfo> get_table_infos(const RelAlgExecutionUnit& ra_exe_un
247247
}
248248

249249
const std::map<int, ChunkMetadata>& Fragmenter_Namespace::FragmentInfo::getChunkMetadataMap() const {
250-
if (resultSet) {
250+
if (resultSet && !synthesizedMetadataIsValid) {
251251
chunkMetadataMap = synthesize_metadata(resultSet);
252-
resultSet = nullptr;
252+
synthesizedMetadataIsValid = true;
253253
}
254254
return chunkMetadataMap;
255255
}
256+
257+
size_t Fragmenter_Namespace::FragmentInfo::getNumTuples() const {
258+
std::unique_ptr<std::lock_guard<std::mutex>> lock;
259+
if (resultSetMutex) {
260+
lock.reset(new std::lock_guard<std::mutex>(*resultSetMutex));
261+
}
262+
if (resultSet && !synthesizedNumTuplesIsValid) {
263+
numTuples = resultSet->rowCount();
264+
synthesizedNumTuplesIsValid = true;
265+
}
266+
return numTuples;
267+
}

QueryEngine/JoinHashTable.cpp

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -174,7 +174,7 @@ std::pair<const int8_t*, size_t> JoinHashTable::getColumnFragment(
174174
effective_mem_lvl,
175175
effective_mem_lvl == Data_Namespace::CPU_LEVEL ? 0 : device_id);
176176
}
177-
return {col_buff, fragment.numTuples};
177+
return {col_buff, fragment.getNumTuples()};
178178
}
179179

180180
std::pair<const int8_t*, size_t> JoinHashTable::getAllColumnFragments(

0 commit comments

Comments
 (0)