| 95 | } |
| 96 | |
| 97 | void WindowStep::transformPipeline(QueryPipeline & pipeline, const BuildQueryPipelineSettings & s) |
| 98 | { |
| 99 | auto enable_windows_parallel = s.context->getSettingsRef().enable_windows_parallel; |
| 100 | if (need_sort && !window_description.full_sort_description.empty()) |
| 101 | { |
| 102 | if (enable_windows_parallel) |
| 103 | scatterByPartitionIfNeeded(pipeline); |
| 104 | |
| 105 | // finish sorting |
| 106 | |
| 107 | DataStream input_stream = input_streams[0]; |
| 108 | if (!prefix_description.empty() && !enable_windows_parallel) |
| 109 | { |
| 110 | bool need_finish_sorting = (prefix_description.size() < window_description.full_sort_description.size()); |
| 111 | if (pipeline.getNumStreams() > 1) |
| 112 | { |
| 113 | auto transform = std::make_shared<MergingSortedTransform>( |
| 114 | pipeline.getHeader(), pipeline.getNumStreams(), prefix_description, s.context->getSettingsRef().max_block_size, 0); |
| 115 | |
| 116 | pipeline.addTransform(std::move(transform)); |
| 117 | } |
| 118 | |
| 119 | if (need_finish_sorting) |
| 120 | { |
| 121 | pipeline.addSimpleTransform([&](const Block & header, QueryPipeline::StreamType stream_type) -> ProcessorPtr { |
| 122 | if (stream_type != QueryPipeline::StreamType::Main) |
| 123 | return nullptr; |
| 124 | |
| 125 | return std::make_shared<PartialSortingTransform>(header, window_description.full_sort_description, 0); |
| 126 | }); |
| 127 | |
| 128 | /// NOTE limits are not applied to the size of temporary sets in FinishSortingTransform |
| 129 | pipeline.addSimpleTransform([&](const Block & header) -> ProcessorPtr { |
| 130 | return std::make_shared<FinishSortingTransform>( |
| 131 | header, |
| 132 | prefix_description, |
| 133 | window_description.full_sort_description, |
| 134 | s.context->getSettingsRef().max_block_size, |
| 135 | 0); |
| 136 | }); |
| 137 | } |
| 138 | } |
| 139 | else |
| 140 | { |
| 141 | PartialSortingStep partial_sorting_step{input_stream, window_description.full_sort_description, 0}; |
| 142 | partial_sorting_step.transformPipeline(pipeline, s); |
| 143 | |
| 144 | MergeSortingStep merge_sorting_step{input_stream, window_description.full_sort_description, 0}; |
| 145 | merge_sorting_step.transformPipeline(pipeline, s); |
| 146 | } |
| 147 | if (!enable_windows_parallel || window_description.partition_by.empty()) |
| 148 | { |
| 149 | MergingSortedStep merging_sorted_step{ |
| 150 | input_stream, window_description.full_sort_description, s.context->getSettingsRef().max_block_size, 0}; |
| 151 | merging_sorted_step.transformPipeline(pipeline, s); |
| 152 | } |
| 153 | } |
| 154 |
nothing calls this directly
no test coverage detected