MCPcopy Create free account
hub / github.com/ClickHouse/ClickHouse / initializePipeline

Method initializePipeline

src/Processors/QueryPlan/GatherReceiveStep.cpp:22–51  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

20{
21
22void GatherReceiveStep::initializePipeline(QueryPipelineBuilder & pipeline, const BuildQueryPipelineSettings & settings)
23{
24 Pipes pipes;
25
26 /// Read from all buckets
27 for (size_t i = 0; i < num_buckets; ++i)
28 {
29 pipes.push_back(Pipe(settings.exchange_lookup->createSource(output_header, ExchangeStreamId(exchange_id, i, 0))));
30 }
31
32 pipeline.init(Pipe::unitePipes(std::move(pipes)));
33
34 if (maintain_sort_description && pipeline.getNumStreams() > 1)
35 {
36 pipeline.addTransform(
37 std::make_shared<MergingSortedTransform>(
38 output_header,
39 num_buckets,
40 *maintain_sort_description,
41 /* merge_block_size_rows */ DEFAULT_BLOCK_SIZE,
42 /* merge_block_size_bytes */ 0,
43 /* max_dynamic_subcolumns */ std::nullopt,
44 SortingQueueStrategy::Batch,
45 /* limit */ 0,
46 /* always_read_till_end */ false,
47 /* rows_sources_write_buf */ nullptr,
48 /* filter_column_name */ std::nullopt,
49 /* blocks_are_granules_size */ false));
50 }
51}
52
53void GatherReceiveStep::serialize(Serialization & ctx) const
54{

Callers

nothing calls this directly

Calls 7

ExchangeStreamIdClass · 0.85
PipeClass · 0.70
push_backMethod · 0.45
createSourceMethod · 0.45
initMethod · 0.45
getNumStreamsMethod · 0.45
addTransformMethod · 0.45

Tested by

no test coverage detected