| 134 | } |
| 135 | |
| 136 | TemporaryBlockStreamHolder SortedBlocksWriter::flush(const BlocksList & blocks) const |
| 137 | { |
| 138 | Pipes pipes; |
| 139 | pipes.reserve(blocks.size()); |
| 140 | for (const auto & block : blocks) |
| 141 | if (auto num_rows = block.rows()) |
| 142 | pipes.emplace_back(std::make_shared<SourceFromSingleChunk>(std::make_shared<const Block>(block.cloneEmpty()), Chunk(block.getColumns(), num_rows))); |
| 143 | |
| 144 | if (pipes.empty()) |
| 145 | throw Exception(ErrorCodes::LOGICAL_ERROR, "Empty block"); |
| 146 | |
| 147 | QueryPipelineBuilder pipeline; |
| 148 | pipeline.init(Pipe::unitePipes(std::move(pipes))); |
| 149 | |
| 150 | if (pipeline.getNumStreams() > 1) |
| 151 | { |
| 152 | auto transform = std::make_shared<MergingSortedTransform>( |
| 153 | pipeline.getSharedHeader(), |
| 154 | pipeline.getNumStreams(), |
| 155 | sort_description, |
| 156 | rows_in_block, |
| 157 | /*max_block_size_bytes=*/0, |
| 158 | /*max_dynamic_subcolumns=*/std::nullopt, |
| 159 | SortingQueueStrategy::Default); |
| 160 | |
| 161 | pipeline.addTransform(std::move(transform)); |
| 162 | } |
| 163 | |
| 164 | return flushToFile(tmp_data, sample_block, std::move(pipeline)); |
| 165 | } |
| 166 | |
| 167 | class TemporaryFileLazySource final : public ISource |
| 168 | { |
nothing calls this directly
no test coverage detected