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]