| 224 | } |
| 225 | |
| 226 | MergeTextIndexesTask::MergeTextIndexesTask( |
| 227 | std::vector<TextIndexSegment> segments_, |
| 228 | MergeTreeMutableDataPartPtr new_data_part_, |
| 229 | size_t num_rows_, |
| 230 | MergeTreeIndexPtr index_ptr_, |
| 231 | std::shared_ptr<MergedPartOffsets> merged_part_offsets_, |
| 232 | const MergeTreeReaderSettings & reader_settings_, |
| 233 | const MergeTreeWriterSettings & writer_settings_) |
| 234 | : segments(std::move(segments_)) |
| 235 | , new_data_part(std::move(new_data_part_)) |
| 236 | , num_rows(num_rows_) |
| 237 | , index_ptr(std::move(index_ptr_)) |
| 238 | , merged_part_offsets(std::move(merged_part_offsets_)) |
| 239 | , writer_settings(writer_settings_) |
| 240 | , step_time_ms((*new_data_part->storage.getSettings())[MergeTreeSetting::background_task_preferred_step_execution_time_ms].totalMilliseconds()) |
| 241 | , postings_serialization(createPostingsSerialization(*index_ptr)) |
| 242 | { |
| 243 | cursors.resize(segments.size()); |
| 244 | inputs.resize(segments.size()); |
| 245 | input_streams.resize(segments.size()); |
| 246 | |
| 247 | output_tokens = ColumnString::create(); |
| 248 | params = typeid_cast<const MergeTreeIndexText &>(*index_ptr).getParams(); |
| 249 | sparse_index_tokens = ColumnString::create(); |
| 250 | sparse_index_offsets = ColumnUInt64::create(); |
| 251 | |
| 252 | std::tie(output_streams, output_streams_holders) = makeOutputStreams( |
| 253 | index_ptr->getSubstreams(), |
| 254 | index_ptr->getFileName(), |
| 255 | new_data_part->getDataPartStoragePtr(), |
| 256 | new_data_part->default_codec, |
| 257 | new_data_part->getMarksFileExtension(), |
| 258 | writer_settings); |
| 259 | |
| 260 | auto substreams = index_ptr->getSubstreams(); |
| 261 | |
| 262 | for (size_t i = 0; i < segments.size(); ++i) |
| 263 | { |
| 264 | for (const auto & substream : substreams) |
| 265 | { |
| 266 | auto stream = makeTextIndexInputStream( |
| 267 | segments[i].part_storage, |
| 268 | segments[i].index_file_name + substream.suffix, |
| 269 | substream.extension, |
| 270 | MergeTreeIndexReader::patchSettings(reader_settings_, substream.type)); |
| 271 | |
| 272 | input_streams[i][substream.type] = stream.get(); |
| 273 | input_streams_holders.emplace_back(std::move(stream)); |
| 274 | } |
| 275 | } |
| 276 | |
| 277 | /// Resolve each source part's codec from its own header. |
| 278 | source_postings_serializations.reserve(segments.size()); |
| 279 | |
| 280 | for (size_t i = 0; i < segments.size(); ++i) |
| 281 | { |
| 282 | auto * stream = input_streams[i].at(MergeTreeIndexSubstream::Type::Regular); |
| 283 | source_postings_serializations.emplace_back(createSourcePostingsSerialization(*stream)); |
nothing calls this directly
no test coverage detected