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

Method getOutputWithSpill

bolt/exec/GroupingSet.cpp:1179–1243  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

1177 BOLT_CHECK(!hasSpilled());
1178
1179 if (table_ == nullptr) {
1180 return;
1181 }
1182
1183 auto* spillConf = const_cast<common::SpillConfig*>(spillConfig_);
1184 operatorCtx_->adjustSpillCompressionKind(spillConf);
1185 spillConf->needSetNextEqual = bypassProbeHT_;
1186 auto* rows = table_->rows();
1187 BOLT_CHECK(pool_.trackUsage());
1188 spiller_ = std::make_unique<Spiller>(
1189 Spiller::Type::kAggregateOutput, rows, makeSpillType(), spillConfig_);
1190 spiller_->setSpillConfig(spillConfig_);
1191 spiller_->spill(rowIterator);
1192 adjustBypassHashTable(spillConf);
1193 table_->clear();
1194}
1195
1196bool GroupingSet::getOutputWithSpill(
1197 int32_t maxOutputRows,
1198 int32_t maxOutputBytes,
1199 const RowVectorPtr& result) {
1200 if (merge_ == nullptr && rowBasedSpillMerge_ == nullptr) {
1201 LOG(INFO) << operatorCtx_->toString() << " prepare merge, files number: "
1202 << spiller_->state().numFinishedFiles(0)
1203 << " memory: " << pool_.currentBytes()
1204 << ", reserved: " << pool_.reservedBytes();
1205 BOLT_CHECK_NULL(mergeRows_);
1206 BOLT_CHECK(mergeArgs_.empty());
1207
1208 if (!isDistinct()) {
1209 mergeArgs_.resize(1);
1210 std::vector<TypePtr> keyTypes;
1211 for (auto& hasher : table_->hashers()) {
1212 keyTypes.push_back(hasher->type());
1213 }
1214
1215 mergeRows_ = std::make_unique<RowContainer>(
1216 keyTypes,
1217 !ignoreNullKeys_,
1218 accumulators(false),
1219 std::vector<TypePtr>(),
1220 false,
1221 false,
1222 true,
1223 false,
1224 false /*useListRowIndex*/,
1225 &pool_,
1226 table_->rows()->stringAllocatorShared());
1227
1228 initializeAggregates(aggregates_, *mergeRows_, false);
1229 }
1230
1231 BOLT_CHECK_EQ(table_->rows()->numRows(), 0);
1232
1233 auto spillPartition = spiller_->finishSpill();
1234 if (spillPartition.numFiles() == 0) {
1235 BOLT_CHECK(spiller_->type() == Spiller::Type::kAggregateOutput);
1236 BOLT_CHECK(spiller_->filledZeroRows());

Callers

nothing calls this directly

Calls 15

initializeAggregatesFunction · 0.85
numFilesMethod · 0.80
filledZeroRowsMethod · 0.80
createOrderedReaderMethod · 0.80
getJITenabledForSpillMethod · 0.80
maxPartitionsMethod · 0.80
toStringMethod · 0.45
numFinishedFilesMethod · 0.45
stateMethod · 0.45
currentBytesMethod · 0.45
reservedBytesMethod · 0.45

Tested by

no test coverage detected