| 856 | } |
| 857 | |
| 858 | static void addMergingFinal( |
| 859 | Pipe & pipe, |
| 860 | size_t num_output_streams, |
| 861 | const SortDescription & sort_description, |
| 862 | MergeTreeMetaBase::MergingParams merging_params, |
| 863 | Names partition_key_columns, |
| 864 | size_t max_block_size) |
| 865 | { |
| 866 | const auto & header = pipe.getHeader(); |
| 867 | size_t num_outputs = pipe.numOutputPorts(); |
| 868 | |
| 869 | auto get_merging_processor = [&]() -> MergingTransformPtr |
| 870 | { |
| 871 | switch (merging_params.mode) |
| 872 | { |
| 873 | case MergeTreeMetaBase::MergingParams::Ordinary: |
| 874 | { |
| 875 | return std::make_shared<MergingSortedTransform>(header, num_outputs, |
| 876 | sort_description, max_block_size); |
| 877 | } |
| 878 | |
| 879 | case MergeTreeMetaBase::MergingParams::Collapsing: |
| 880 | return std::make_shared<CollapsingSortedTransform>(header, num_outputs, |
| 881 | sort_description, merging_params.sign_column, true, max_block_size); |
| 882 | |
| 883 | case MergeTreeMetaBase::MergingParams::Summing: |
| 884 | return std::make_shared<SummingSortedTransform>(header, num_outputs, |
| 885 | sort_description, merging_params.columns_to_sum, partition_key_columns, max_block_size); |
| 886 | |
| 887 | case MergeTreeMetaBase::MergingParams::Aggregating: |
| 888 | return std::make_shared<AggregatingSortedTransform>(header, num_outputs, |
| 889 | sort_description, max_block_size); |
| 890 | |
| 891 | case MergeTreeMetaBase::MergingParams::Replacing: |
| 892 | return std::make_shared<ReplacingSortedTransform>(header, num_outputs, |
| 893 | sort_description, merging_params.version_column, max_block_size); |
| 894 | |
| 895 | case MergeTreeMetaBase::MergingParams::VersionedCollapsing: |
| 896 | return std::make_shared<VersionedCollapsingTransform>(header, num_outputs, |
| 897 | sort_description, merging_params.sign_column, max_block_size); |
| 898 | |
| 899 | case MergeTreeMetaBase::MergingParams::Graphite: |
| 900 | throw Exception("GraphiteMergeTree doesn't support FINAL", ErrorCodes::LOGICAL_ERROR); |
| 901 | } |
| 902 | |
| 903 | __builtin_unreachable(); |
| 904 | }; |
| 905 | |
| 906 | if (num_output_streams <= 1 || sort_description.empty()) |
| 907 | { |
| 908 | pipe.addTransform(get_merging_processor()); |
| 909 | return; |
| 910 | } |
| 911 | |
| 912 | ColumnNumbers key_columns; |
| 913 | key_columns.reserve(sort_description.size()); |
| 914 | |
| 915 | for (const auto & desc : sort_description) |
no test coverage detected