| 15 | { |
| 16 | |
| 17 | void BroadcastReceiveStep::initializePipeline(QueryPipelineBuilder & pipeline, const BuildQueryPipelineSettings & settings) |
| 18 | { |
| 19 | const String bucket_id = settings.parameter_lookup->getParameter("bucket_id").safeGet<String>(); |
| 20 | |
| 21 | VectorWithMemoryTracking<std::unique_ptr<QueryPipelineBuilder>> pipelines; |
| 22 | |
| 23 | /// Read all shards |
| 24 | for (const String & shard_id : source_shards) |
| 25 | { |
| 26 | std::unique_ptr<QueryPipelineBuilder> pipeline_ptr = std::make_unique<QueryPipelineBuilder>(); |
| 27 | pipeline_ptr->init(Pipe(settings.exchange_lookup->createSource(output_header, ExchangeStreamId(exchange_id, shard_id, bucket_id)))); |
| 28 | pipelines.emplace_back(std::move(pipeline_ptr)); |
| 29 | } |
| 30 | |
| 31 | pipeline = QueryPipelineBuilder::unitePipelines(std::move(pipelines), 0, &processors); |
| 32 | } |
| 33 | |
| 34 | void BroadcastReceiveStep::serialize(Serialization & ctx) const |
| 35 | { |
nothing calls this directly
no test coverage detected