yujun777 opened a new issue, #68679:
URL: https://github.com/apache/doris/issues/68679

   ### Search before asking
   
   - [X] I searched in the [issues](https://github.com/apache/doris/issues) and 
found nothing similar.
   
   ### Version
   
   master (the two-phase overwrite in `InsertOverwriteTableCommand` and its 
publication via `ReplacePartitionOperationLog`, and the remote/cloud commit 
paths as of 2026-09).
   
   ### What's Wrong?
   
   An `INSERT OVERWRITE` is two-phase: the rows are committed into temporary 
partitions first (`InsertIntoTableCommand` hangs the table stream offsets on 
that transaction; `DatabaseTransactionMgr.updateCatalogAfterCommitted` applies 
them on commit and replays them), and a later swap publishes them 
(`Env.replaceTempPartition` / `ReplacePartitionOperationLog`). Two durable 
facts are therefore split across two journal records, and three remaining paths 
can leave the offset (and the rows) advanced while nothing is published, with 
no error:
   
   1. **The FE stops between the two halves** (crash, or a master switch, where 
`Env.transferToMaster` calls `InsertOverwriteManager.allTaskFail()` and drops 
every in-flight task's temporary partitions) or **the swap itself fails** 
(`replacePartition` throwing). The committed rows are gone and the consumed 
offset stays advanced, so a re-run reads from the advanced position and that 
range is missing.
   2. **A remote commit whose reply is lost.** 
`RemoteOlapInsertExecutor.onComplete` commits on the owning FE and then waits; 
if `masterCallWithRetry` exhausts its retries after the owner committed, this 
FE sees an error, `onFail`'s abort cannot take a committed remote transaction 
back, and the overwrite rolls the temporary partitions back, dropping rows that 
exist.
   3. **A cloud transaction that is maybe-committed.** FoundationDB can apply a 
commit and still report an error; after the finite meta-service retries this 
reaches the FE as `KV_TXN_COMMIT_ERR` and 
`commitAndPublishTransactionWithRetry` throws, so the overwrite rolls back 
temporary partitions whose rows and stream offsets the meta service may already 
hold.
   
   The same shape shows up where the decision is the owning FE's: for a remote 
target, a cancellation that arrives while `replacePartitionsImpl` waits for the 
table write lock is not conveyed by the replacement RPC, so an overwrite whose 
insert committed nothing can still publish an empty result and report success. 
(The local equivalents of the cancellation cases, and the local error-response 
case, are fixed by #68662.)
   
   Impact: silent partial data for `INSERT OVERWRITE t SELECT * FROM 
stream(...)`, and for an IVM/MTMV partition refresh, which resets the stream 
offsets of exactly the partitions it replaces -- the partition keeps its old 
rows, the change it read is consumed, and the following refresh reports 
success. Trace issue for the IVM side: #65418.
   
   ### What You Expected?
   
   Once a load's rows (or stream offsets) are committed, the overwrite must 
publish them: `committed` and `published` must not be able to disagree, 
whatever happens to the FE, the RPC or the meta service in between. Where the 
commit outcome is genuinely unknown, the overwrite must not answer with a 
rollback that can drop durable rows.
   
   ### How to Reproduce?
   
   1. Local, crash or swap failure: run `INSERT OVERWRITE t SELECT * FROM 
stream('db', 'src')` (or an IVM partition refresh) and stop the FE (or inject a 
failure in the swap) between the insert's commit and the swap; the temporary 
partitions are dropped by `allTaskFail`/`taskFail` while the offsets stay 
advanced.
   2. Remote: overwrite a remote Doris table and drop the commit reply 
(`masterCallWithRetry` exhausting retries) after the owning FE committed; the 
temporary partitions are rolled back on this side.
   3. Cloud: overwrite with a commit that the meta service applies but reports 
as failed (FDB maybe-committed -> `KV_TXN_COMMIT_ERR`).
   
   The local cancellation and publication-timeout variants are pinned by the 
regression suite added in #68662 
(`insert_overwrite_p0/test_insert_overwrite_cancel`), and the local 
`replacePartition`-failure and crash windows are pinned by 
`mtmv_p0/ivm/test_ivm_overwrite_failure_between_the_halves`.
   
   ### Anything Else?
   
   Fix directions, in increasing order of cost:
   
   - Carry the swap as a committed action of the insert transaction, so one 
journal record decides both the offset advance and the publication (the offsets 
already ride `TransactionState`; a "committed catalog action" for the partition 
swap would ride the same record, and the publish/finish path already re-drives 
committed transactions after a restart).
   - Or move the consumption position out of the insert transaction and into 
the swap's record (`ReplacePartitionOperationLog`), which keeps the transaction 
path untouched but weakens the read-range exclusivity `checkStreamOffset` 
provides during the window, and needs a separate answer for cloud read state.
   - For remote and cloud specifically, the owner side (the owning FE, the meta 
service) has to be able to answer "did that commit happen"; until then, an 
uncertain commit has no correct local resolution -- dropping the partitions 
loses durable rows and keeping them leaves unpublished ones that the next 
`allTaskFail` drops anyway.
   
   ### Are you willing to submit PR?
   
   - [X] Yes I am willing to submit a PR!
   
   ### Code of Conduct
   
   - [X] I agree to follow this project's [Code of 
Conduct](https://github.com/apache/doris/blob/master/CODE_OF_CONDUCT.md)
   


-- 
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