MCPcopy Create free account
hub / github.com/bytedance/bolt / next

Method next

bolt/exec/LocalPartition.cpp:237–263  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

235}
236
237BlockingReason LocalExchangeQueue::next(
238 ContinueFuture* future,
239 memory::MemoryPool* pool,
240 RowVectorPtr* data) {
241 std::vector<ContinuePromise> memoryPromises;
242 auto blockingReason = queue_.withWLock([&](auto& queue) {
243 *data = nullptr;
244 if (queue.empty()) {
245 if (isFinishedLocked(queue)) {
246 return BlockingReason::kNotBlocked;
247 }
248
249 consumerPromises_.emplace_back("LocalExchangeQueue::next");
250 *future = consumerPromises_.back().getSemiFuture();
251
252 return BlockingReason::kWaitForProducer;
253 }
254 *data = queue.front();
255 queue.pop();
256 memoryPromises =
257 memoryManager_->decreaseMemoryUsage((*data)->estimateFlatSize());
258
259 return BlockingReason::kNotBlocked;
260 });
261 notify(memoryPromises);
262 return blockingReason;
263}
264
265bool LocalExchangeQueue::isFinishedLocked(
266 const std::queue<RowVectorPtr>& queue) const {

Callers 1

getOutputMethod · 0.45

Calls 6

notifyFunction · 0.85
backMethod · 0.80
emptyMethod · 0.45
popMethod · 0.45
decreaseMemoryUsageMethod · 0.45
estimateFlatSizeMethod · 0.45

Tested by

no test coverage detected