github-actions[bot] commented on code in PR #68032:
URL: https://github.com/apache/doris/pull/68032#discussion_r4230579122


##########
be/src/exec/spill/spill_file_manager.cpp:
##########
@@ -153,14 +199,16 @@ void SpillFileManager::stop() {
     // The GC thread may observe the stop latch before processing a recently 
queued failed deletion.
     // Retry the pending directories after the thread exits; later failures 
get one final retry in
     // the destructor.
-    _retry_pending_query_spill_directories();
+    _retry_pending_spill_directories();
 }
 
 Status SpillFileManager::init() {
     LOG(INFO) << "init spill stream manager";
     RETURN_IF_ERROR(_init_spill_store_map());
+    _remote_upload_budget =
+            
std::make_shared<SpillRemoteUploadBudget>(config::spill_s3_max_inflight_upload_bytes);
 
-    for (const auto& [path, store] : _spill_store_map) {
+    for (auto* store : _local_stores) {

Review Comment:
   [P1] Recover completed spill objects before reporting a fresh total after BE 
restart. `_persisted_bytes` is process-local and incremented only when a live 
`SpillFile` closes a part; this initialization scans local spill roots but 
never reconciles the S3 objects. After an abrupt crash leaves completed objects 
behind, the restarted BE returns zero and every FE stamps that as a fresh 
billing total until bucket expiration. The lifecycle-rule decision covers 
deletion, while the PR promises accounting after restart. Reconcile only this 
BE/instance's residual objects or use a durable accounting source, and test 
write, hard restart, then poll.



##########
be/src/cloud/cloud_storage_engine.cpp:
##########
@@ -452,12 +452,18 @@ Status 
CloudStorageEngine::start_bg_threads(std::shared_ptr<WorkloadGroup> wg_sp
 void CloudStorageEngine::sync_storage_vault() {
     cloud::StorageVaultInfos vault_infos;
     bool enable_storage_vault = false;
+    std::string default_vault_id;
 
-    auto st = _meta_mgr->get_storage_vault_info(&vault_infos, 
&enable_storage_vault);
+    auto st = _meta_mgr->get_storage_vault_info(&vault_infos, 
&enable_storage_vault,
+                                                &default_vault_id);
     if (!st.ok()) {
         LOG(WARNING) << "failed to get storage vault info. err=" << st;
         return;
     }
+    {
+        std::lock_guard lock(_latest_fs_mtx);
+        _default_vault_id = std::move(default_vault_id);

Review Comment:
   [P2] Publish the new default vault only after its filesystem is available. 
This assigns the new ID before the loop creates and registers a newly selected 
vault. A spill that starts in that interval resolves the new ID, gets no 
filesystem, and fails with `storage vault ... not found`. If filesystem 
creation fails, the refresh only logs a warning and leaves the unusable ID 
published until another successful refresh. Install the filesystem first, then 
switch the default, and cover concurrent selection of a newly added vault.



##########
fe/fe-core/src/main/java/org/apache/doris/plugin/audit/AuditLoader.java:
##########
@@ -269,32 +273,31 @@ public synchronized void loadIfNecessary(boolean force) {
 
         if (auditLogBuffer.length() != 0 && (force || auditLogBuffer.length() 
>= GlobalVariable.auditPluginMaxBatchBytes
                 || currentTime - lastLoadTimeAuditLog >= 
GlobalVariable.auditPluginMaxBatchInternalSec * 1000)) {
-            // begin to load
+            // Keep the batch until a successful load. In a follower-first 
upgrade the table can
+            // still have its old schema, and a rejected request must not 
erase audit rows.
             try {
-                String token = "";
-                try {
-                    // Acquire token from master
-                    token = 
Env.getCurrentEnv().getTokenManager().acquireToken();
-                } catch (Exception e) {
-                    LOG.warn("Failed to get auth token: {}", e);
-                    discardLogNum += auditLogNum;
-                    return;
+                if (batchLabel == null) {
+                    batchLabel = streamLoader.genLabel();
                 }
-                AuditStreamLoader.LoadResponse response = 
streamLoader.loadBatch(auditLogBuffer, token);
-                if (LOG.isDebugEnabled()) {
-                    LOG.debug("audit loader response: {}", response);
-                }
-            } catch (Exception e) {
-                if (LOG.isDebugEnabled()) {
-                    LOG.debug("encounter exception when putting current audit 
batch, discard current batch", e);
+                String token = 
Env.getCurrentEnv().getTokenManager().acquireToken();
+                AuditStreamLoader.LoadResponse response = 
streamLoader.loadBatch(auditLogBuffer, token, batchLabel);
+                if (!response.succeeded(auditLogNum)) {
+                    failedLoad = true;
+                    if (response.rejectedOrIncomplete(auditLogNum)) {

Review Comment:
   [P2] Preserve the batch label when `Fail` can follow an uncertain commit. 
`rejectedOrIncomplete()` treats every BE `Status: Fail` as proof that nothing 
committed, but the BE also emits `Fail` when its FE commit RPC loses the reply 
after `loadTxnCommitImpl` committed. This branch then creates a new label and 
retries the same rows into the duplicate-key audit table, producing duplicate 
audit events. Classify a definitive pre-commit rejection separately from an 
uncertain commit outcome, and test a committed transaction whose commit 
acknowledgment is lost.



##########
be/src/exec/spill/spill_file_manager.cpp:
##########
@@ -431,10 +622,10 @@ void SpillFileManager::gc(int32_t max_work_time_ms) {
             LOG(INFO) << msg;
         }
     }};
-    _retry_pending_query_spill_directories();
-    for (const auto& [path, store_dir] : _spill_store_map) {
+    _retry_pending_spill_directories();

Review Comment:
   [P2] Bound remote delete retries before entering the local GC pass. A failed 
S3 `delete_directory` queues each spill-file prefix, but this call drains 
**every** prefix synchronously before `spill_gc_work_time_ms` is checked. When 
the object store is unavailable, one pass can serially wait for many network 
timeouts; `stop()` then joins that pass and repeats the entire backlog, 
delaying BE shutdown. Process a bounded number/time slice, retain the rest, and 
make shutdown skip or bound remote retries.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to