| 63 | namespace |
| 64 | { |
| 65 | void transformToSingleBlockSources(Pipes & inputs) |
| 66 | { |
| 67 | size_t inputs_size = inputs.size(); |
| 68 | for (size_t i = 0; i < inputs_size; ++i) |
| 69 | { |
| 70 | auto && input = inputs[i]; |
| 71 | QueryPipeline input_pipeline(std::move(input)); |
| 72 | PullingPipelineExecutor input_pipeline_executor(input_pipeline); |
| 73 | |
| 74 | auto header = input_pipeline_executor.getHeader(); |
| 75 | auto result_block = header.cloneEmpty(); |
| 76 | |
| 77 | size_t result_block_columns = result_block.columns(); |
| 78 | |
| 79 | Block result; |
| 80 | while (input_pipeline_executor.pull(result)) |
| 81 | { |
| 82 | for (size_t result_block_index = 0; result_block_index < result_block_columns; ++result_block_index) |
| 83 | { |
| 84 | auto & block_column = result.safeGetByPosition(result_block_index); |
| 85 | auto & result_block_column = result_block.safeGetByPosition(result_block_index); |
| 86 | |
| 87 | auto mutable_column = IColumn::mutate(std::move(result_block_column.column)); |
| 88 | mutable_column->insertRangeFrom(*block_column.column, 0, block_column.column->size()); |
| 89 | result_block_column.column = std::move(mutable_column); |
| 90 | } |
| 91 | } |
| 92 | |
| 93 | auto source = std::make_shared<SourceFromSingleChunk>(std::make_shared<const Block>(std::move(result_block))); |
| 94 | inputs[i] = Pipe(std::move(source)); |
| 95 | } |
| 96 | } |
| 97 | } |
| 98 | |
| 99 | StorageExecutable::StorageExecutable( |
no test coverage detected