MCPcopy Create free account
hub / github.com/ByConity/ByConity / transformPipeline

Method transformPipeline

src/QueryPlan/WindowStep.cpp:97–168  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

95}
96
97void 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

Callers

nothing calls this directly

Calls 8

getNumStreamsMethod · 0.80
emptyMethod · 0.45
sizeMethod · 0.45
getHeaderMethod · 0.45
addTransformMethod · 0.45
addSimpleTransformMethod · 0.45
resizeMethod · 0.45

Tested by

no test coverage detected