morningman commented on PR #68371:
URL: https://github.com/apache/doris/pull/68371#issuecomment-5844972250
## P1 — Publish the full set of report state the waiter reads, or release
the latch at the end of report handling
The fix reorders two scalar fields, but the same root cause still applies to
everything else the report handler writes after the latch has already been
released.
`LoadLoadingTask` wakes up inside `Coordinator.updateStatus()` (via
`cancelInternal` → `cancelLatch` → `fragmentsDoneLatch.countDownToZero()`), and
immediately snapshots:
```java
attachment = new BrokerLoadingTaskAttachment(signature,
curCoordinator.getLoadCounters(), // also racing
curCoordinator.getTrackingUrl(),
curCoordinator.getFirstErrorMsg(),
TabletCommitInfo.fromThrift(curCoordinator.getCommitInfos()),
ErrorTabletInfo.fromThrift(curCoordinator.getErrorTabletInfos()...),
status);
curCoordinator.getErrorTabletInfos().clear();
```
Three reasons this matters:
1. **The report that loses the URL carries the counters too.** On the BE
side `params.load_counters.emplace(...)` is set unconditionally, immediately
next to `if (!req.load_error_url.empty()) params.__set_tracking_url(...)`
(`be/src/exec/pipeline/pipeline_fragment_context.cpp:2481-2490`). So
`DPP_ABNORMAL_ALL` / `UNSELECTED_ROWS` are lost with the same probability as
the URL, and `SHOW LOAD` shows zero filtered rows for the failed task.
2. **Those getters return live collections**: `getLoadCounters()` →
`loadCounters` (`Coordinator.java:475`), `getCommitInfos()` → `commitInfos`
(`:519`), `getErrorTabletInfos()` → `errorTabletInfos` (`:523`). The report
thread does `put` / `addAll` on plain `HashMap` / `ArrayList` while the waiter
reads them — the worst case is not merely a stale empty value but a partially
built list, a `ConcurrentModificationException`, or `clear()` interleaving with
`addAll`.
3. The trailing `getErrorTabletInfos().clear()` widens that window further.
The Nereids path has the same shape (`acceptFinalReport` runs after
`updateStatusIfOk`), but `LoadContext`'s accessors are `synchronized` with
`CopyOnWriteArrayList` and `ImmutableMap.copyOf` (`LoadContext.java:43-111`),
so the collection hazard is much smaller there; the remaining risk is
`AbstractInsertExecutor` / `RemoteOlapInsertExecutor` reading incomplete values
after the wait.
**Minimal fix** — move these fields along with the diagnostics. All four are
self-contained under `lock` and do not depend on `reportTxnId`, which is
resolved later:
```java
Status status = new Status(params.status);
// Publish everything a waiting load task snapshots after join() returns,
before an error
// status can release the completion latch from inside updateStatus().
if (params.isSetTrackingUrl()) { trackingUrl = params.getTrackingUrl(); }
if (params.isSetFirstErrorMsg()) { firstErrorMsg =
params.getFirstErrorMsg(); }
if (params.isSetDeltaUrls() && deltaUrls != null) {
updateDeltas(params.getDeltaUrls()); }
if (params.isSetLoadCounters() && loadCounters != null) {
updateLoadCounters(params.getLoadCounters()); }
if (params.isSetCommitInfos()) {
updateCommitInfos(params.getCommitInfos()); }
if (params.isSetErrorTabletInfos()) {
updateErrorTabletInfos(params.getErrorTabletInfos()); }
if (!status.ok()) { ... updateStatus(status); }
```
(Do not move the Hive/Iceberg/MC commit-data feeding — that one does depend
on `reportTxnId`.)
**Structural fix** — strengthen the invariant instead: do not release the
latch from `cancelInternal` on the report-driven failure path; release it when
report handling finishes. `cancelInternal` keeps notifying remote fragments for
cancellation, and the external paths (`Coordinator.cancel()`, timeouts, BE-down
detection) still need to release immediately. This makes the invariant hold for
fields added later as well.
Either way, if the remaining fields are deliberately out of scope, say so in
the PR description. The current non-goal only mentions "reports arriving after
task teardown"; it does not mention the sibling fields carried by the very same
report.
## P2 — Document the hook contract
The hook's correctness depends entirely on two rules: it must run before the
error status releases waiters, and implementations must be idempotent (because
`doProcessReportExecStatus` applies the same fields again). A one-line comment
does not carry that:
```java
/**
* Publishes report diagnostics that must be visible to a waiting load task
before an error
* status releases its latch. Called before {@code updateStatusIfOk} on the
failure path, and
* therefore also invoked for reports whose final aggregation runs later in
* {@link #doProcessReportExecStatus}; implementations must be idempotent.
*/
protected void publishReportDiagnosticsBeforeStatus(TReportExecStatusParams
params) {}
```
## P2 — Explain why `updateLoadDiagnostics` runs twice on the failure path
On failure it is invoked once by the hook and once by `acceptFinalReport`.
It is harmless (`LoadContext.updateTrackingUrl` / `updateFirstErrorMsg` are
plain assignments), but nothing in the code says it is intentional, so the next
reader will assume a leftover. A short comment at the `acceptFinalReport` call
site would settle it, e.g. `// Re-applied here for the success path; already
published by the pre-status hook on failure.`
--
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]