| 1561 | } |
| 1562 | |
| 1563 | static void addMergingFinal( |
| 1564 | Pipe & pipe, |
| 1565 | const SortDescription & sort_description, |
| 1566 | MergeTreeData::MergingParams merging_params, |
| 1567 | const StorageMetadataPtr & metadata_snapshot, |
| 1568 | size_t max_block_size_rows, |
| 1569 | bool enable_vertical_final) |
| 1570 | { |
| 1571 | auto header = pipe.getSharedHeader(); |
| 1572 | size_t num_outputs = pipe.numOutputPorts(); |
| 1573 | |
| 1574 | auto now = time(nullptr); |
| 1575 | |
| 1576 | auto get_merging_processor = [&]() -> MergingTransformPtr |
| 1577 | { |
| 1578 | switch (merging_params.mode) |
| 1579 | { |
| 1580 | case MergeTreeData::MergingParams::Ordinary: |
| 1581 | return std::make_shared<MergingSortedTransform>(header, num_outputs, |
| 1582 | sort_description, max_block_size_rows, /*max_block_size_bytes=*/0, /*max_dynamic_subcolumns*/std::nullopt, SortingQueueStrategy::Batch); |
| 1583 | |
| 1584 | case MergeTreeData::MergingParams::Collapsing: |
| 1585 | return std::make_shared<CollapsingSortedTransform>(header, num_outputs, |
| 1586 | sort_description, merging_params.sign_column, true, max_block_size_rows, /*max_block_size_bytes=*/0, /*max_dynamic_subcolumns*/std::nullopt); |
| 1587 | |
| 1588 | case MergeTreeData::MergingParams::Summing: { |
| 1589 | auto required_columns = metadata_snapshot->getPartitionKey().expression->getRequiredColumns(); |
| 1590 | required_columns.append_range(metadata_snapshot->getSortingKey().expression->getRequiredColumns()); |
| 1591 | return std::make_shared<SummingSortedTransform>(header, num_outputs, |
| 1592 | sort_description, merging_params.columns_to_sum, required_columns, max_block_size_rows, /*max_block_size_bytes=*/0, /*max_dynamic_subcolumns*/std::nullopt, merging_params.allow_tuple_element_aggregation); |
| 1593 | } |
| 1594 | |
| 1595 | case MergeTreeData::MergingParams::Aggregating: |
| 1596 | return std::make_shared<AggregatingSortedTransform>(header, num_outputs, |
| 1597 | sort_description, max_block_size_rows, /*max_block_size_bytes=*/0, /*max_dynamic_subcolumns*/std::nullopt, merging_params.allow_tuple_element_aggregation); |
| 1598 | |
| 1599 | case MergeTreeData::MergingParams::Replacing: |
| 1600 | return std::make_shared<ReplacingSortedTransform>(header, num_outputs, |
| 1601 | sort_description, merging_params.is_deleted_column, merging_params.version_column, max_block_size_rows, /*max_block_size_bytes=*/0, /*max_dynamic_subcolumns*/std::nullopt, /*out_row_sources_buf_*/ nullptr, /*use_average_block_sizes*/ false, /*cleanup*/ !merging_params.is_deleted_column.empty(), enable_vertical_final); |
| 1602 | |
| 1603 | |
| 1604 | case MergeTreeData::MergingParams::VersionedCollapsing: |
| 1605 | return std::make_shared<VersionedCollapsingTransform>(header, num_outputs, |
| 1606 | sort_description, merging_params.sign_column, max_block_size_rows, /*max_block_size_bytes=*/0, /*max_dynamic_subcolumns*/std::nullopt); |
| 1607 | |
| 1608 | case MergeTreeData::MergingParams::Graphite: |
| 1609 | return std::make_shared<GraphiteRollupSortedTransform>(header, num_outputs, |
| 1610 | sort_description, max_block_size_rows, /*max_block_size_bytes=*/0, /*max_dynamic_subcolumns*/std::nullopt, merging_params.graphite_params, now); |
| 1611 | |
| 1612 | case MergeTreeData::MergingParams::Coalescing: |
| 1613 | { |
| 1614 | auto required_columns = metadata_snapshot->getPartitionKey().expression->getRequiredColumns(); |
| 1615 | required_columns.append_range(metadata_snapshot->getSortingKey().expression->getRequiredColumns()); |
| 1616 | return std::make_shared<CoalescingSortedTransform>(header, num_outputs, |
| 1617 | sort_description, merging_params.columns_to_sum, required_columns, max_block_size_rows, /*max_block_size_bytes=*/0, /*max_dynamic_subcolumns*/std::nullopt, merging_params.allow_tuple_element_aggregation); |
| 1618 | } |
| 1619 | } |
| 1620 | }; |
no test coverage detected