MCPcopy Create free account
hub / github.com/apache/arrow / NextLineItemBatch

Method NextLineItemBatch

cpp/src/arrow/acero/tpch_node.cc:1388–1465  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

1386 }
1387
1388 Result<std::optional<ExecBatch>> NextLineItemBatch(size_t thread_index) {
1389 ThreadLocalData& tld = thread_local_data_[thread_index];
1390 ExecBatch queued;
1391 bool from_queue = false;
1392 {
1393 std::lock_guard<std::mutex> lock(lineitem_output_queue_mutex_);
1394 if (!lineitem_output_queue_.empty()) {
1395 queued = std::move(lineitem_output_queue_.front());
1396 lineitem_output_queue_.pop();
1397 from_queue = true;
1398 }
1399 }
1400 tld.first_batch_offset = 0;
1401 if (from_queue) {
1402 ARROW_DCHECK(queued.length <= batch_size_);
1403 tld.first_batch_offset = queued.length;
1404 if (queued.length == batch_size_) return queued;
1405 }
1406 {
1407 std::lock_guard<std::mutex> lock(orders_output_queue_mutex_);
1408 if (orders_rows_generated_ == orders_rows_to_generate_) {
1409 if (from_queue) return queued;
1410 return std::nullopt;
1411 }
1412
1413 tld.orderkey_start = orders_rows_generated_;
1414 tld.orders_to_generate =
1415 std::min(batch_size_, orders_rows_to_generate_ - orders_rows_generated_);
1416 orders_rows_generated_ += tld.orders_to_generate;
1417 orders_batches_generated_.fetch_add(1);
1418 RETURN_NOT_OK(GenerateRowCounts(thread_index));
1419 lineitem_batches_generated_.fetch_add(
1420 static_cast<int64_t>(tld.lineitem.size() - from_queue));
1421 ARROW_DCHECK(orders_rows_generated_ <= orders_rows_to_generate_);
1422 }
1423 tld.orders.resize(ORDERS::kNumCols);
1424 std::fill(tld.orders.begin(), tld.orders.end(), Datum());
1425 tld.generated_lineitem.reset();
1426 if (from_queue) {
1427 for (size_t i = 0; i < lineitem_cols_.size(); i++)
1428 if (tld.lineitem[0][lineitem_cols_[i]].kind() == Datum::NONE)
1429 tld.lineitem[0][lineitem_cols_[i]] = std::move(queued[i]);
1430 }
1431
1432 for (int col : orders_cols_) RETURN_NOT_OK(kOrdersGenerators[col](thread_index));
1433 for (int col : lineitem_cols_) RETURN_NOT_OK(kLineitemGenerators[col](thread_index));
1434
1435 if (!orders_cols_.empty()) {
1436 std::vector<Datum> orders_result(orders_cols_.size());
1437 for (size_t i = 0; i < orders_cols_.size(); i++) {
1438 int col_idx = orders_cols_[i];
1439 orders_result[i] = tld.orders[col_idx];
1440 }
1441 ARROW_ASSIGN_OR_RAISE(ExecBatch orders_batch,
1442 ExecBatch::Make(std::move(orders_result)));
1443 {
1444 std::lock_guard<std::mutex> lock(orders_output_queue_mutex_);
1445 orders_output_queue_.emplace(std::move(orders_batch));

Callers 1

ProduceCallbackMethod · 0.80

Calls 9

resizeMethod · 0.80
emplace_backMethod · 0.80
DatumClass · 0.50
emptyMethod · 0.45
sizeMethod · 0.45
beginMethod · 0.45
endMethod · 0.45
resetMethod · 0.45
kindMethod · 0.45

Tested by

no test coverage detected