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]
