From 0d3b87e65725f2e7a0861d0072471a06edd2e105 Mon Sep 17 00:00:00 2001 From: bobhan1 Date: Wed, 29 Jul 2026 15:06:50 +0800 Subject: [PATCH] [fix](cloud) Prevent urgent load from preempting schema change lock (#66082) Related PR: #53086 Problem Summary: Urgent MOW load lock requests currently replace any existing delete bitmap update lock after the force-take threshold. If a heavy schema change owns the lock, the load can invalidate the schema change and make its later commit fail. This change keeps active schema change locks protected in both v1 and v2 lock implementations. Urgent loads continue to preempt compaction and load locks, and stale schema change locks still follow the existing expiration cleanup path. The unit-test matrix directly verifies persisted lock ownership and v2 job-key cleanup for active schema change, expired schema change, compaction, and ordinary load owners. An isolated cloud Docker case pauses schema change after lock acquisition, limits the urgent load to two immediate lock attempts, and verifies that the load fails with the lock conflict while schema change retains ownership. ### Release note Urgent MOW loads no longer preempt active schema change delete bitmap locks. --- cloud/src/meta-service/meta_service.cpp | 8 +- cloud/test/meta_service_test.cpp | 118 +++++++++++++++++- ..._mow_load_preserves_schema_change_lock.out | 4 + ...w_load_preserves_schema_change_lock.groovy | 101 +++++++++++++++ 4 files changed, 226 insertions(+), 5 deletions(-) create mode 100644 regression-test/data/fault_injection_p0/cloud/test_cloud_mow_load_preserves_schema_change_lock.out create mode 100644 regression-test/suites/fault_injection_p0/cloud/test_cloud_mow_load_preserves_schema_change_lock.groovy diff --git a/cloud/src/meta-service/meta_service.cpp b/cloud/src/meta-service/meta_service.cpp index 2fe48980f6f193..525e6391e0546c 100644 --- a/cloud/src/meta-service/meta_service.cpp +++ b/cloud/src/meta-service/meta_service.cpp @@ -4576,6 +4576,10 @@ static bool put_delete_bitmap_update_lock_key(MetaServiceCode& code, std::string return true; } +static bool can_load_force_take_delete_bitmap_lock(int64_t current_lock_id) { + return current_lock_id != SCHEMA_CHANGE_DELETE_BITMAP_LOCK_ID; +} + bool MetaServiceImpl::get_mow_tablet_stats_and_meta(MetaServiceCode& code, std::string& msg, const GetDeleteBitmapUpdateLockRequest* request, GetDeleteBitmapUpdateLockResponse* response, @@ -4950,7 +4954,7 @@ void MetaServiceImpl::get_delete_bitmap_update_lock_v2( msg = "failed to parse DeleteBitmapUpdateLockPB"; return; } - if (urgent) { + if (urgent && can_load_force_take_delete_bitmap_lock(lock_info.lock_id())) { // since currently only the FE Master initiates the lock request for import tasks, // and it does so in a single-threaded manner, there is no need to check the lock id here DCHECK(request->lock_id() > 0); @@ -5204,7 +5208,7 @@ void MetaServiceImpl::get_delete_bitmap_update_lock_v1( msg = "failed to parse DeleteBitmapUpdateLockPB"; return; } - if (urgent) { + if (urgent && can_load_force_take_delete_bitmap_lock(lock_info.lock_id())) { // since currently only the FE Master initiates the lock request for import tasks, // and it does so in a single-threaded manner, there is no need to check the lock id here DCHECK(request->lock_id() > 0); diff --git a/cloud/test/meta_service_test.cpp b/cloud/test/meta_service_test.cpp index 9b04e98a2e0010..be3425c3736323 100644 --- a/cloud/test/meta_service_test.cpp +++ b/cloud/test/meta_service_test.cpp @@ -4924,6 +4924,28 @@ void remove_delete_bitmap_lock(MetaServiceProxy* meta_service, int64_t table_id) ASSERT_EQ(txn->commit(), TxnErrorCode::TXN_OK); } +void check_delete_bitmap_lock_id(MetaServiceProxy* meta_service, int64_t table_id, + int64_t expected_lock_id) { + std::unique_ptr txn; + ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), TxnErrorCode::TXN_OK); + std::string lock_key = meta_delete_bitmap_update_lock_key({"test_instance", table_id, -1}); + std::string lock_val; + ASSERT_EQ(txn->get(lock_key, &lock_val), TxnErrorCode::TXN_OK); + DeleteBitmapUpdateLockPB lock_info; + ASSERT_TRUE(lock_info.ParseFromString(lock_val)); + EXPECT_EQ(lock_info.lock_id(), expected_lock_id); +} + +void check_mow_tablet_job_key(MetaServiceProxy* meta_service, int64_t table_id, int64_t initiator, + bool expected_exists) { + std::unique_ptr txn; + ASSERT_EQ(meta_service->txn_kv()->create_txn(&txn), TxnErrorCode::TXN_OK); + std::string job_key = mow_tablet_job_key({"test_instance", table_id, initiator}); + std::string job_val; + EXPECT_EQ(txn->get(job_key, &job_val), + expected_exists ? TxnErrorCode::TXN_OK : TxnErrorCode::TXN_KEY_NOT_FOUND); +} + void testGetDeleteBitmapUpdateLock(int lock_version, int job_lock_id) { config::delete_bitmap_lock_v2_white_list = lock_version == 1 ? "" : "*"; auto meta_service = get_meta_service(); @@ -5140,7 +5162,33 @@ void testGetDeleteBitmapUpdateLock(int lock_version, int job_lock_id) { nullptr); ASSERT_EQ(remove_res.status().code(), MetaServiceCode::OK); - // case 11: lock by schema change but expired, compaction get lock but txn commit conflict, do fast retry + // case 11: urgent load can force take compaction lock but not schema change lock + req.set_lock_id(job_lock_id); + req.set_initiator(100); + req.set_expiration(100); + meta_service->get_delete_bitmap_update_lock( + reinterpret_cast<::google::protobuf::RpcController*>(&cntl), &req, &res, nullptr); + ASSERT_EQ(res.status().code(), MetaServiceCode::OK); + + req.set_lock_id(888); + req.set_initiator(-1); + req.set_expiration(60); + req.set_urgent(true); + meta_service->get_delete_bitmap_update_lock( + reinterpret_cast<::google::protobuf::RpcController*>(&cntl), &req, &res, nullptr); + ASSERT_EQ(res.status().code(), job_lock_id == SCHEMA_CHANGE_DELETE_BITMAP_LOCK_ID + ? MetaServiceCode::LOCK_CONFLICT + : MetaServiceCode::OK); + req.set_urgent(false); + remove_req.set_lock_id(job_lock_id == SCHEMA_CHANGE_DELETE_BITMAP_LOCK_ID ? job_lock_id : 888); + remove_req.set_initiator(job_lock_id == SCHEMA_CHANGE_DELETE_BITMAP_LOCK_ID ? 100 : -1); + meta_service->remove_delete_bitmap_update_lock( + reinterpret_cast<::google::protobuf::RpcController*>(&cntl), &remove_req, &remove_res, + nullptr); + ASSERT_EQ(remove_res.status().code(), MetaServiceCode::OK); + remove_delete_bitmap_lock(meta_service.get(), table_id); + + // case 12: lock by schema change but expired, compaction get lock but txn commit conflict, do fast retry sp->set_call_back("get_delete_bitmap_update_lock:commit:conflict", [&](auto&& args) { auto* first_retry = try_any_cast(args[0]); auto lock_id = (try_any_cast(args[1]))->lock_id(); @@ -5164,7 +5212,7 @@ void testGetDeleteBitmapUpdateLock(int lock_version, int job_lock_id) { nullptr); ASSERT_EQ(remove_res.status().code(), MetaServiceCode::OK); - // case 12: lock by load but expired, compaction get lock but txn commit conflict, do fast retry + // case 13: lock by load but expired, compaction get lock but txn commit conflict, do fast retry req.set_lock_id(300); req.set_initiator(-1); req.set_expiration(1); @@ -5181,7 +5229,7 @@ void testGetDeleteBitmapUpdateLock(int lock_version, int job_lock_id) { ASSERT_EQ(res.status().code(), MetaServiceCode::OK); remove_delete_bitmap_lock(meta_service.get(), table_id); - // case 13: lock key does not exist, compaction get lock but txn commit conflict, do fast retry + // case 14: lock key does not exist, compaction get lock but txn commit conflict, do fast retry meta_service->get_delete_bitmap_update_lock( reinterpret_cast<::google::protobuf::RpcController*>(&cntl), &req, &res, nullptr); ASSERT_EQ(res.status().code(), MetaServiceCode::OK); @@ -5199,6 +5247,70 @@ TEST(MetaServiceTest, GetDeleteBitmapUpdateLock) { testGetDeleteBitmapUpdateLock(1, SCHEMA_CHANGE_DELETE_BITMAP_LOCK_ID); } +void testUrgentLoadDeleteBitmapLock(int lock_version) { + config::delete_bitmap_lock_v2_white_list = lock_version == 1 ? "" : "*"; + auto meta_service = get_meta_service(); + int64_t table_id = 90 + lock_version; + remove_delete_bitmap_lock(meta_service.get(), table_id); + + brpc::Controller cntl; + GetDeleteBitmapUpdateLockRequest req; + GetDeleteBitmapUpdateLockResponse res; + req.set_cloud_unique_id("test_cloud_unique_id"); + req.set_table_id(table_id); + req.add_partition_ids(123); + + auto get_lock = [&](int64_t lock_id, int64_t initiator, int64_t expiration, bool urgent) { + req.set_lock_id(lock_id); + req.set_initiator(initiator); + req.set_expiration(expiration); + req.set_urgent(urgent); + res.Clear(); + meta_service->get_delete_bitmap_update_lock( + reinterpret_cast<::google::protobuf::RpcController*>(&cntl), &req, &res, nullptr); + return res.status().code(); + }; + + // An urgent load must preserve an active schema change lock. + ASSERT_EQ(get_lock(SCHEMA_CHANGE_DELETE_BITMAP_LOCK_ID, 100, 100, false), MetaServiceCode::OK); + ASSERT_EQ(get_lock(888, -1, 60, true), MetaServiceCode::LOCK_CONFLICT); + check_delete_bitmap_lock_id(meta_service.get(), table_id, SCHEMA_CHANGE_DELETE_BITMAP_LOCK_ID); + if (lock_version == 2) { + check_mow_tablet_job_key(meta_service.get(), table_id, 100, true); + } + remove_delete_bitmap_lock(meta_service.get(), table_id); + + // Expired schema change locks still follow the ordinary stale-lock cleanup path. + ASSERT_EQ(get_lock(SCHEMA_CHANGE_DELETE_BITMAP_LOCK_ID, 101, 1, false), MetaServiceCode::OK); + sleep(2); + ASSERT_EQ(get_lock(888, -1, 60, true), MetaServiceCode::OK); + check_delete_bitmap_lock_id(meta_service.get(), table_id, 888); + if (lock_version == 2) { + check_mow_tablet_job_key(meta_service.get(), table_id, 101, false); + } + remove_delete_bitmap_lock(meta_service.get(), table_id); + + // The existing force-take behavior for compaction locks is unchanged. + ASSERT_EQ(get_lock(COMPACTION_DELETE_BITMAP_LOCK_ID, 102, 100, false), MetaServiceCode::OK); + ASSERT_EQ(get_lock(888, -1, 60, true), MetaServiceCode::OK); + check_delete_bitmap_lock_id(meta_service.get(), table_id, 888); + if (lock_version == 2) { + check_mow_tablet_job_key(meta_service.get(), table_id, 102, false); + } + remove_delete_bitmap_lock(meta_service.get(), table_id); + + // The existing force-take behavior for another load lock is unchanged. + ASSERT_EQ(get_lock(777, -1, 100, false), MetaServiceCode::OK); + ASSERT_EQ(get_lock(888, -1, 60, true), MetaServiceCode::OK); + check_delete_bitmap_lock_id(meta_service.get(), table_id, 888); + remove_delete_bitmap_lock(meta_service.get(), table_id); +} + +TEST(MetaServiceTest, UrgentLoadDeleteBitmapLock) { + testUrgentLoadDeleteBitmapLock(2); + testUrgentLoadDeleteBitmapLock(1); +} + TEST(MetaServiceTest, GetDeleteBitmapUpdateLockNoReadStats) { auto meta_service = get_meta_service(); diff --git a/regression-test/data/fault_injection_p0/cloud/test_cloud_mow_load_preserves_schema_change_lock.out b/regression-test/data/fault_injection_p0/cloud/test_cloud_mow_load_preserves_schema_change_lock.out new file mode 100644 index 00000000000000..f51f1281f0a897 --- /dev/null +++ b/regression-test/data/fault_injection_p0/cloud/test_cloud_mow_load_preserves_schema_change_lock.out @@ -0,0 +1,4 @@ +-- This file is automatically generated. You should know what you did if you want to edit this +-- !final -- +1 10 + diff --git a/regression-test/suites/fault_injection_p0/cloud/test_cloud_mow_load_preserves_schema_change_lock.groovy b/regression-test/suites/fault_injection_p0/cloud/test_cloud_mow_load_preserves_schema_change_lock.groovy new file mode 100644 index 00000000000000..8316bb0f4c8634 --- /dev/null +++ b/regression-test/suites/fault_injection_p0/cloud/test_cloud_mow_load_preserves_schema_change_lock.groovy @@ -0,0 +1,101 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +import org.apache.doris.regression.suite.ClusterOptions + +suite("test_cloud_mow_load_preserves_schema_change_lock", "docker") { + def options = new ClusterOptions() + options.setFeNum(1) + options.setBeNum(1) + options.cloudMode = true + options.enableDebugPoints() + options.feConfigs += [ + "delete_bitmap_lock_expiration_seconds=60", + "enable_mow_load_force_take_ms_lock=true", + "mow_load_force_take_ms_lock_threshold_ms=0", + "meta_service_rpc_retry_times=2", + "enable_schema_change_retry=false" + ] + options.beConfigs += [ + "delete_bitmap_lock_expiration_seconds=60" + ] + + docker(options) { + GetDebugPoint().clearDebugPointsForAllBEs() + + try { + sql "DROP TABLE IF EXISTS test_cloud_mow_load_preserves_schema_change_lock" + sql """ + CREATE TABLE test_cloud_mow_load_preserves_schema_change_lock ( + k INT NOT NULL, + v INT + ) + UNIQUE KEY(k) + DISTRIBUTED BY HASH(k) BUCKETS 1 + PROPERTIES ( + "enable_unique_key_merge_on_write" = "true", + "disable_auto_compaction" = "true", + "light_schema_change" = "false", + "replication_num" = "1" + ) + """ + sql "INSERT INTO test_cloud_mow_load_preserves_schema_change_lock VALUES (1, 10)" + + GetDebugPoint().enableDebugPointForAllBEs( + "CloudSchemaChangeJob::_process_delete_bitmap.inject_sleep", + [percent: "1.0", sleep: "10"]) + def alterFuture = thread { + sql """ + ALTER TABLE test_cloud_mow_load_preserves_schema_change_lock + MODIFY COLUMN v BIGINT + """ + } + // This debug point is after the schema change acquires the delete bitmap lock. + // The debug point runs after schema change acquires the delete bitmap lock. + // The table has only one rowset, so three seconds is ample for reaching it. + sleep(3000) + + test { + sql "INSERT INTO test_cloud_mow_load_preserves_schema_change_lock VALUES (2, 20)" + exception "Failed to get delete bitmap lock due to conflict" + } + + alterFuture.get() + waitForSchemaChangeDone { + sql """ + SHOW ALTER TABLE COLUMN + WHERE TableName = 'test_cloud_mow_load_preserves_schema_change_lock' + ORDER BY CreateTime DESC LIMIT 1 + """ + time 120 + } + def alterResult = sql_return_maparray """ + SHOW ALTER TABLE COLUMN + WHERE TableName = 'test_cloud_mow_load_preserves_schema_change_lock' + ORDER BY CreateTime DESC LIMIT 1 + """ + assertEquals("FINISHED", alterResult[0].State) + + order_qt_final """ + SELECT k, v FROM test_cloud_mow_load_preserves_schema_change_lock + ORDER BY k + """ + } finally { + GetDebugPoint().clearDebugPointsForAllBEs() + } + } +}