| 44 | } |
| 45 | |
| 46 | void CloudMergeTreeMutateTask::executeImpl() |
| 47 | { |
| 48 | auto lock_holder = storage.lockForShare(RWLockImpl::NO_QUERY, storage.getSettings()->lock_acquire_timeout_for_background_operations); |
| 49 | |
| 50 | MergeTreeDataMutator mutate_executor(storage, getContext()->getSettingsRef().background_pool_size); |
| 51 | auto data_parts = mutate_executor.mutatePartsToTemporaryParts(params, *manipulation_entry, getContext(), lock_holder); |
| 52 | |
| 53 | if (isCancelled()) |
| 54 | throw Exception("Merge task " + params.task_id + " is cancelled", ErrorCodes::ABORTED); |
| 55 | |
| 56 | UInt64 peak_memory_usage = 0; |
| 57 | |
| 58 | ManipulationListElement * manipulation_list_element = getManipulationListElement(); |
| 59 | if (manipulation_list_element) |
| 60 | { |
| 61 | peak_memory_usage = manipulation_list_element->getMemoryTracker().getPeak(); |
| 62 | } |
| 63 | |
| 64 | CnchDataWriter cnch_writer(storage, getContext(), ManipulationType::Mutate, params.task_id, |
| 65 | /*consumer_group_*/ {}, /*tpl_*/ {}, /*binlog*/ {}, peak_memory_usage); |
| 66 | auto res = cnch_writer.dumpAndCommitCnchParts(data_parts); |
| 67 | getContext()->getCurrentTransaction()->commitV2(); |
| 68 | if (params.parts_preload_level) |
| 69 | cnch_writer.preload(res.parts); |
| 70 | } |
| 71 | |
| 72 | } |
nothing calls this directly
no test coverage detected