@@ -102,6 +102,67 @@ bvar::Adder<int64_t> g_peer_server_fill_rejected("peer_server_fill_rejected");
102102bvar::LatencyRecorder g_peer_server_fill_latency (" peer_server_fill_latency" );
103103bvar::LatencyRecorder g_cloud_internal_service_get_file_cache_meta_by_tablet_id_latency (
104104 " cloud_internal_service_get_file_cache_meta_by_tablet_id_latency" );
105+ bvar::Adder<int64_t > g_cloud_sync_tablet_meta_requests_total (
106+ " cloud_sync_tablet_meta_requests_total" );
107+ bvar::Adder<int64_t > g_cloud_sync_tablet_meta_synced_total (" cloud_sync_tablet_meta_synced_total" );
108+ bvar::Adder<int64_t > g_cloud_sync_tablet_meta_skipped_total (" cloud_sync_tablet_meta_skipped_total" );
109+ bvar::Adder<int64_t > g_cloud_sync_tablet_meta_failed_total (" cloud_sync_tablet_meta_failed_total" );
110+
111+ namespace {
112+
113+ void submit_sync_tablet_meta (CloudStorageEngine& engine, FifoThreadPool& work_pool,
114+ const PSyncTabletMetaRequest* request,
115+ PSyncTabletMetaResponse* response, google::protobuf::Closure* done) {
116+ auto start_time = std::chrono::steady_clock::now ();
117+ bool ret = work_pool.try_offer ([engine = &engine, request, response, done, start_time]() {
118+ brpc::ClosureGuard closure_guard (done);
119+ LOG (INFO ) << " begin to sync tablet meta, request=" << request->ShortDebugString ();
120+ int64_t synced = 0 ;
121+ int64_t skipped = 0 ;
122+ int64_t failed = 0 ;
123+ g_cloud_sync_tablet_meta_requests_total << 1 ;
124+ for (const auto tablet_id : request->tablet_ids ()) {
125+ auto tablet = engine->tablet_mgr ().get_tablet_if_cached (tablet_id);
126+ if (!tablet) {
127+ ++skipped;
128+ continue ;
129+ }
130+ auto st = tablet->sync_meta ();
131+ if (!st.ok ()) {
132+ ++failed;
133+ LOG (WARNING ) << " failed to sync tablet meta from cloud meta service, tablet="
134+ << tablet_id << " , err=" << st;
135+ continue ;
136+ }
137+ ++synced;
138+ }
139+ g_cloud_sync_tablet_meta_synced_total << synced;
140+ g_cloud_sync_tablet_meta_skipped_total << skipped;
141+ g_cloud_sync_tablet_meta_failed_total << failed;
142+ response->set_synced_tablets (synced);
143+ response->set_skipped_tablets (skipped);
144+ response->set_failed_tablets (failed);
145+ Status::OK ().to_protobuf (response->mutable_status ());
146+ auto cost_ms = std::chrono::duration_cast<std::chrono::milliseconds>(
147+ std::chrono::steady_clock::now () - start_time)
148+ .count ();
149+ LOG (INFO ) << " finish to sync tablet meta, request=" << request->ShortDebugString ()
150+ << " , response=" << response->ShortDebugString () << " , cost_ms=" << cost_ms;
151+ });
152+ if (!ret) {
153+ brpc::ClosureGuard closure_guard (done);
154+ Status::InternalError (" failed to offer sync_tablet_meta request to work pool" )
155+ .to_protobuf (response->mutable_status ());
156+ auto cost_ms = std::chrono::duration_cast<std::chrono::milliseconds>(
157+ std::chrono::steady_clock::now () - start_time)
158+ .count ();
159+ LOG (WARNING ) << " failed to offer sync_tablet_meta request to work pool, request="
160+ << request->ShortDebugString () << " , response=" << response->ShortDebugString ()
161+ << " , cost_ms=" << cost_ms;
162+ }
163+ }
164+
165+ } // namespace
105166
106167// Concurrency guard for server-side S3 pull-through fills.
107168static std::atomic<int32_t > g_active_server_fills {0 };
@@ -114,6 +175,22 @@ CloudInternalServiceImpl::CloudInternalServiceImpl(CloudStorageEngine& engine, E
114175
115176CloudInternalServiceImpl::~CloudInternalServiceImpl () = default ;
116177
178+ void CloudInternalServiceImpl::sync_tablet_meta (google::protobuf::RpcController* controller,
179+ const PSyncTabletMetaRequest* request,
180+ PSyncTabletMetaResponse* response,
181+ google::protobuf::Closure* done) {
182+ submit_sync_tablet_meta (_engine, _light_work_pool, request, response, done);
183+ }
184+
185+ #ifdef BE_TEST
186+ void test_submit_sync_tablet_meta (CloudStorageEngine& engine, FifoThreadPool& work_pool,
187+ const PSyncTabletMetaRequest* request,
188+ PSyncTabletMetaResponse* response,
189+ google::protobuf::Closure* done) {
190+ submit_sync_tablet_meta (engine, work_pool, request, response, done);
191+ }
192+ #endif
193+
117194void CloudInternalServiceImpl::alter_vault_sync (google::protobuf::RpcController* controller,
118195 const doris::PAlterVaultSyncRequest* request,
119196 PAlterVaultSyncResponse* response,
@@ -925,12 +1002,7 @@ bvar::Adder<uint64_t> g_file_cache_warm_up_rowset_wait_for_compaction_num(
9251002bvar::Adder<uint64_t > g_file_cache_warm_up_rowset_wait_for_compaction_timeout_num (
9261003 " file_cache_warm_up_rowset_wait_for_compaction_timeout_num" );
9271004
928- // Per-job windowed metrics for target BE
929- // bvar::Window enforces MAX_SECONDS_LIMIT = 3600, so the longest window is 1h.
930- static constexpr int WINDOW_5M = 300 ;
931- static constexpr int WINDOW_30M = 1800 ;
932- static constexpr int WINDOW_1H = 3600 ;
933-
1005+ // Per-job windowed metrics for target BE (window spans shared via bvar_windowed_adder.h)
9341006MBvarWindowedAdder g_warmup_ed_finish_segment_num (" warmup_ed_finish_segment_num" , {" job_id" },
9351007 {WINDOW_5M , WINDOW_30M , WINDOW_1H }, false );
9361008MBvarWindowedAdder g_warmup_ed_finish_segment_size (" warmup_ed_finish_segment_size" , {" job_id" },
0 commit comments