| 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)); |
no test coverage detected