Skip to content

[feature](cloud) Spill to object storage in cloud mode and report the traffic in SHOW DATA - #68032

Open
mrhhsg wants to merge 23 commits into
apache:masterfrom
mrhhsg:feat/cloud-spill-s3
Open

mrhhsg wants to merge 23 commits into
apache:masterfrom
mrhhsg:feat/cloud-spill-s3

Conversation

@mrhhsg

@mrhhsg mrhhsg commented Sep 15, 2026 •

Copy link
Copy Markdown
Member

What problem does this PR solve?

Issue Number: None

Problem Summary:

Cloud-mode BEs could only spill to local disks. This PR adds spill to the instance's S3 storage vault, selected by be.conf (spill_storage_type = local | s3, mutually exclusive), while keeping the existing spill operators and part-file format. The spill data held in object storage is a billing input, so every FE polls the BEs for the size they currently hold and SHOW DATA shows the sum.

BE

  • RemoteSpillDataDir bound to the storage vault (spill_s3_storage_vault, default: the instance's default vault); objects live under spill/{ip}_{port}/{query_id}/{spill file}/{part}, where ip is the address the BE advertises and port its heartbeat_service_port (the pair FE identifies a BE by), so several BEs on one host never share a directory and a key tells which BE and which query wrote it. The layout assumes that no other BE writing to the vault has the same address and port. The store becomes ready lazily (the vault may not be known at startup). A query registers its directory before its first object is written; once the store is ready, the GC thread lists spill/{ip}_{port}/ once and deletes, one per round, the query directories that no query of this process registered, i.e. the residue of the previous process.
  • SpillFileWriter/SpillFileReader reuse S3FileWriter/S3FileReader. Atomic capacity reservation (spill_s3_storage_limit_bytes) and a submit-time upload budget (spill_s3_max_inflight_upload_bytes, charged with the allocated capacity of every in-flight upload buffer) that blocks writers when too many upload buffers are in flight and refuses cancelled queries; FileWriterOptions gains an upload gate + done callback for this, honoured by S3FileWriter. The spill store is split into SpillDataDir (base), LocalSpillDataDir and RemoteSpillDataDir.
  • Parts are closed synchronously (the wait at a part boundary is negligible against the default 1 GB part, and spilling already runs on an IO thread); failed parts are drained before the budget is reconciled and multipart uploads are aborted.
  • bvars/metrics: spill_remote_{read,write}_bytes, spill_remote_{get,put}_requests and per-second throughput/QPS; profile counters SpillRemote*.
  • Per-query statistics: TQueryStatistics gains spill_write_bytes_to_remote_storage / spill_read_bytes_from_remote_storage, filled by the BE.
  • get_be_resource (the brpc the FE already uses for admission control) additionally returns the spill bytes the BE currently holds in object storage (PGlobalResourceUsage.remote_spill_bytes). Upload traffic (bytes and PutObject/UploadPart requests) stays available as BE metrics/bvars.
  • RuntimeState releases the operator local states before its object pool. The pool owns the profiles whose counters the local states hold; a spill writer left open (cancelling a query makes a remote close fail, and the hash join sink stops closing writers at the first error) closes its last part in its destructor and used a freed timer.

Meta-service

  • Recycler task recycle_expired_spill_objects removes objects under spill/ older than spill_objects_expire_time_second (default 7 days) as a safety net for crashed BEs that never come back (S3/MOCK accessors). The keys do not name the instance, so in a vault shared by several instances the sweep also removes their expired spill objects, and the snapshot-enabled deleted-instance path runs the same expiration-based sweep instead of dropping a prefix; without snapshots the whole vault is deleted anyway. Incomplete multipart uploads of a crashed BE are invisible to prefix listings; an AbortIncompleteMultipartUpload bucket lifecycle rule is the intended safety net for those.

FE

  • SHOW DATA PROPERTIES("entire_warehouse"="true") lists the spill bytes currently held in object storage by the alive BEs as a row of its own, __remote_spill__ (a user database name has to start with a letter, so it cannot collide), and counts them in the total row. The row is added in cloud mode for the whole warehouse only, not when db_names restricts the listing, and the column layout is unchanged. Every FE polls all alive BEs with get_be_resource once per tablet_stat_update_interval_second (CloudTabletStatMgr, like the tablet sizes) and SHOW DATA reads the cached sum; a poll with any failed BE keeps the previous value, and a value that is missing or older than cloud_spill_stats_max_age_second (new FE config, default 300 s) is reported as an error instead of being shown as 0, because the value is a billing input. Objects left behind by a BE that is gone are not counted until its restart or the recycler removes them; meta-service holds no spill statistics.
  • The audit log (fe.audit.log and __internal_schema.audit_log) gains SpillWriteBytesToRemoteStorage / SpillReadBytesFromRemoteStorage (columns spill_write_bytes_to_remote_storage / spill_read_bytes_from_remote_storage), filled from the BE-reported query statistics like the local spill bytes; an existing audit table gets the columns through the usual schema check at startup.

Design notes and the local review records are kept outside the repository.

Release note

Cloud mode: spill can be written to the S3 storage vault (spill_storage_type = s3 in be.conf); SHOW DATA PROPERTIES("entire_warehouse"="true") shows the spill data currently held in object storage in a __remote_spill__ row that is included in the total; the audit log records the spill bytes written to and read from object storage per query.

Check List (For Author)

  • Test

    • Regression test
    • Unit Test
    • Manual test (add detailed scripts or steps below)
    • No need to test or manual test. Explain why:
      • This is a refactor/code format and no logic has been changed.
      • Previous test can cover this change.
      • No code files have been changed.
      • Other reason
  • Behavior changed:

    • No.
    • Yes. New be.conf options (spill_storage_type, spill_s3_storage_vault, spill_s3_storage_limit_bytes, spill_s3_max_inflight_upload_bytes), new meta-service config spill_objects_expire_time_second, new fe.conf option cloud_spill_stats_max_age_second, new __remote_spill__ row in SHOW DATA ... entire_warehouse (included in the total), two new audit log columns.
  • Does this need documentation?

    • No.
    • Yes.

Check List (For Reviewer who merge this PR)

  • Confirm the release note
  • Confirm test cases
  • Confirm document
  • Add branch pick label

https://claude.ai/code/session_01Jrdwwwh8bZSnVwoCwykHse
https://claude.ai/code/session_01CA1ot1uEwvpQzCJUhmdH32

@hello-stephen

Copy link
Copy Markdown
Contributor

Thank you for your contribution to Apache Doris.
Don't know what should be done next? See How to process your PR.

Please clearly describe your PR:

  1. What problem was fixed (it's best to include specific error reporting information). How it was fixed.
  2. Which behaviors were modified. What was the previous behavior, what is it now, why was it modified, and what possible impacts might there be.
  3. What features were added. Why was this function added?
  4. Which code was refactored and why was this part of the code refactored?
  5. Which functions were optimized and what is the difference before and after the optimization?

@mrhhsg

mrhhsg commented Sep 15, 2026

Copy link
Copy Markdown
Member Author

/review

@github-actions github-actions Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Static verdict: REQUEST_CHANGES. The S3-backed spill path and warehouse-level SHOW DATA plumbing are present, but 16 blocking correctness and operability issues remain. The highest risks are cross-BE deletion caused by a shared identity, non-durable and non-idempotent lifetime statistics, cleanup paths that can delete live data or leak billable storage, request counters that do not measure physical provider traffic, and unbounded recovery/metadata scans.

Critical checkpoints (Part 1.3):

  1. Goal and proof: The goal is remote S3 spill with capacity control, request/byte accounting, persisted warehouse totals, cleanup, and SHOW DATA exposure. Happy-path BE/cloud/FE/regression tests were added, but the goal is not correctly achieved because identity, durability, cleanup, retry, and accounting failures remain.
  2. Scope and focus: The 47-file change is cohesive around that feature, and the cross-layer breadth is necessary, but the accepted defects mean it is not yet review-ready.
  3. Concurrency: I traced the spill-GC thread, writer/upload pools, non-blocking close, upload-budget mutex/condition variable, atomic counters, heartbeat mutation, and FDB transactions. Reachable upload callbacks are balanced and I found no separate deadlock, but release builds mask an underflow invariant (MMF-10), and transaction serialization does not fence stale generations (MMF-9).
  4. Lifecycle and SIOF: No cross-TU static-initialization dependency was found. Process boot, query/file teardown, startup cleanup, graceful/abrupt shutdown, vault binding, and multipart ownership were traced; MMF-4, MMF-6, MMF-7, MMF-8, MMF-11, MMF-12, and MMF-13 show lifecycle failures.
  5. Configuration: New mutable storage, upload-budget, and recycler controls are observed online, but negative storage limits disable the cap (MMF-3), non-positive TTLs can select live objects (MMF-7), and runtime accepts/persists an upload budget that startup rejects (MMF-15).
  6. Compatibility: The protobuf additions are additive, SHOW DATA row width is consistent, non-cloud mode avoids the RPC, and the repository's supported latest-MS/older-FE direction remains compatible. The FE-first/older-MS hypothesis was dismissed under that explicit contract.
  7. Parallel paths: Local and S3 writer paths, PutObject/UploadPart, data/footer reads, synchronous/non-blocking close, and hooks with no stats sink were checked. Local behavior remains balanced, but both read and write request accounting miss physical retries/failures (MMF-5, MMF-16).
  8. Special conditions: Remote-only selection, default-vault lookup, positive-limit gates, age cutoffs, and release underflow handling were checked. The non-positive and defensive-continue conditions are not safe (MMF-3, MMF-7, MMF-10, MMF-15).
  9. Test coverage: Ordinary round trips, rotation, cleanup retries, capacity concurrency, budget blocking/cancellation, several upload failures, meta-service aggregation, and SHOW DATA shape are covered. Missing negative coverage includes two live BEs sharing an identity, reordered/stale reports, crash before report, live identity migration, real expiry/liveness, abort failure/restart residue, paginated cleanup, churn cardinality, invalid dynamic updates, and SDK/GET retries.
  10. Test results: The changed expected SHOW DATA output is internally consistent with the five-column metadata. Per the review-only contract, I did not run builds or tests; author/CI results were not independently verified here.
  11. Observability: Profiles, workload counters, BE metrics, logs, and meta-service bvars were added, but the report RPC omits its FDB read metrics (MMF-1), physical S3 attempts are undercounted (MMF-5/MMF-16), and the persisted total itself is not reliable across crashes or identity transitions.
  12. Transactions and persistence: The FDB read-modify-write is atomic per attempt, but unequal boot IDs are not an ordering fence (MMF-9), process counters lose unreported crash suffixes (MMF-6), and retired identity records are never compacted (MMF-14). No FE EditLog change is involved.
  13. Data writes and crash behavior: S3 multipart creation/upload/abort, visible-object deletion, capacity reservation, and stats writes were traced. Crashes can lose totals and multipart ownership, deletion failure releases live capacity, and cleanup can delete another live BE or a long-running query's data (MMF-4/MMF-6/MMF-7/MMF-8/MMF-11).
  14. FE/BE variable passing: No new scattered FE-to-BE session variable path was introduced; the new MS/FE protobuf and five-column mapping are consistent. The existing heartbeat cloud identity is nevertheless shared and mutable in ways incompatible with its new use as storage/stats identity (MMF-4/MMF-13).
  15. Performance: Startup recursively materializes the entire spill namespace outside its work budget (MMF-12), and SHOW DATA range-reads every historical backend identity without pagination or compaction (MMF-14).
  16. Other issues: Three capped normal/risk convergence rounds found no additional distinct issue after the 16 inline findings. Two hypotheses were dismissed with code evidence: dynamic default-vault rebinding is not an explicit contract, and FE-first upgrade is outside the documented MS-first rollout direction.

Focus response: no additional user-provided review focus was supplied; full-scope review was performed.

Comment thread cloud/src/meta-service/meta_service.cpp Outdated
Comment thread fe/fe-core/src/main/java/org/apache/doris/cloud/rpc/MetaServiceProxy.java Outdated
Comment thread be/src/common/config.cpp
Comment thread be/src/exec/spill/spill_file_manager.cpp Outdated
Comment thread be/src/exec/spill/spill_file_reader.cpp Outdated
Comment thread be/src/exec/spill/spill_file_manager.cpp Outdated
Comment thread be/src/cloud/cloud_meta_mgr.cpp Outdated
Comment thread cloud/src/meta-service/meta_service.cpp Outdated
Comment thread be/src/exec/spill/spill_file_manager.cpp
Comment thread be/src/io/fs/s3_file_writer.cpp
@mrhhsg

mrhhsg commented Sep 16, 2026

Copy link
Copy Markdown
Member Author

Thanks for the review. Pushed a follow-up commit addressing it; per finding:

Fixed

  • meta_service.cpp RPC_PREPROCESS: now (report_spill_stats, get, put) with the get counter/byte bvars.
  • MetaServiceProxy.getSpillStats: goes through executeWithMetrics (standard retry / reconnect / MS_TOO_BUSY handling).
  • spill_s3_storage_limit_bytes: validator rejects negative values; spill_s3_max_inflight_upload_bytes: validator rejects <= 0, and the startup check for < 2 * s3_write_buffer_size now matches the runtime path (warning + degraded one-buffer-at-a-time mode) instead of being fatal.
  • Shared identity (spill_file_manager.cpp / cloud_meta_mgr.cpp): you are right that cloud_unique_id is shared by every BE added in one statement. Both the object root and the stats key now use the FE-assigned backend_id (spill/{backend_id}/data/{boot_id}/..., stats/{instance}/spill/{backend_id}); SpillStatsPB.backend_id added, cloud_unique_id kept as an informational field. A live BE whose id changes (DROP + ADD) keeps the id it was bound with and logs a warning; the new id takes effect at restart.
  • Stale reports (meta_service.cpp fold): a report whose boot_id is older than the recorded one is rejected (INVALID_ARGUMENT) instead of being folded again; a same-boot report never rolls totals back (max). Test covers old-after-new delivery.
  • Startup cleanup (spill_file_manager.cpp list): each boot now writes a small marker spill/{backend_id}/boots/{boot_id} (refreshed daily); cleanup lists only that directory and deletes one old generation per GC round, so the listing is bounded by the number of boot generations, not by the number of spill objects.
  • spill_remote_upload_budget.cpp: DORIS_CHECK_GE, no clamp.
  • recycler.cpp: spill_objects_expire_time_second <= 0 skips the task with a warning (test added); config comment documents the "must exceed the longest query" contract.
  • Proto comment: reporting is documented as best-effort (a crash loses at most one reporting interval; graceful shutdown reports after all tasks are done).

Kept as designed (happy to discuss)

  • Request counters (spill_file_reader.cpp, s3_file_writer.cpp): they count logical SDK calls, not physical retries. Attempt-level accounting needs a per-request hook inside the SDK retry strategy (S3CustomRetryStrategy is process-global with no request context); the byte counters, which are the primary billing input, are exact. Documented as the counting unit; can be a follow-up.
  • Durable per-upload checkpoint (spill_file_manager.cpp counters): writing to meta-service on every upload is not worth the cost; the loss bound is one reporting interval and only on crash.
  • Capacity release on failed delete (spill_file.cpp): same behaviour as the local spill path; retained objects are reclaimed by the query-directory retry and the recycler.
  • Incomplete multipart uploads (spill_file_writer.cpp abort): failed uploads are aborted; residue of a crashed process is invisible to any prefix listing, and neither the BE object client nor the recycler has ListMultipartUploads. The bucket-level AbortIncompleteMultipartUpload rule is the intended safety net (now stated in the config comments).
  • Stats cardinality (get_spill_stats): with backend_id the family grows with "backends ever created" rather than with restarts. Folding retired backends on DROP BACKEND is a reasonable follow-up.
  • Recycler liveness fence: the TTL is documented as "larger than the longest query the cluster allows"; a query holding spill for more than 7 days is outside what this PR targets.

@mrhhsg
mrhhsg force-pushed the feat/cloud-spill-s3 branch from add6dfd to 02ba809 Compare September 17, 2026 02:02
@mrhhsg

mrhhsg commented Sep 17, 2026

Copy link
Copy Markdown
Member Author

Local pipeline review — ✅ PASS

schema: doris-repo-review/v1
status: PASS
pr: apache/doris#68032
commit: 5d93bb00ab769d313f5868a764aafc3aa77c9865
base: 3c489e297d95191c26e7259a9939c2fabb3510cb
reviewed_at: 2026-09-17T15:55+08:00
reviewer: mrhhsg
model: gpt-6-astra
effort: xhigh
findings: {blocker: 0, major: 0, minor: 3, nit: 1}
rounds: 1
converged: true

Notes for maintainers

  • 本次为只读静态评审,覆盖 59/59 个变更文件;未运行构建、UT、回归或真实对象存储验证。
  • 既有 inline threads 已读取并去重;本收据不表示这些历史问题已解决,也不替代人工 Apache approval。

Reviewed locally with the doris-repo-review pipeline. Repository policy may accept this receipt for the matching commit; it is not a human Apache approval.

@mrhhsg
mrhhsg force-pushed the feat/cloud-spill-s3 branch 2 times, most recently from 0389435 to 94124f3 Compare September 19, 2026 14:27
… traffic in SHOW DATA

Cloud-mode BEs can now spill to the instance's S3 storage vault instead of
local disks, selected by be.conf (`spill_storage_type = local | s3`). The
existing spill operators and part-file format are unchanged; a remote
`SpillDataDir` binds the vault file system lazily (no meta-service sync at
startup) and `SpillFileWriter`/`SpillFileReader` reuse `S3FileWriter` and
`S3FileReader`.

BE:
- Remote spill store keyed by `spill/{cloud_unique_id}/{boot_id}/...`; other
  boot generations are cleaned by the GC thread once the store is ready.
- Atomic capacity reservation (`spill_s3_storage_limit_bytes`) and a
  submit-time upload budget (`spill_s3_max_inflight_upload_bytes`) that
  backpressures writers; `FileWriterOptions` gains an upload gate and a
  done callback for that, honoured by `S3FileWriter`.
- Failed parts are drained before the budget is reconciled; multipart
  uploads are aborted on failure.
- bvars/metrics for remote spill bytes, requests and throughput; profile
  counters for remote spill.
- Since-boot upload totals are reported to meta-service about once a
  minute and once at shutdown after all tasks are done (bounded retries).

Meta-service:
- `report_spill_stats` / `get_spill_stats`: one record per BE
  (`stats/{instance}/spill/{cloud_unique_id}`); a new boot folds the
  previous process into `prior_boots_*`, so the record count is bounded by
  the number of BEs. Records are removed with the instance.
- Recycler task `recycle_expired_spill_objects` deletes spill objects older
  than `spill_objects_expire_time_second` (S3/MOCK accessors).

FE:
- `SHOW DATA PROPERTIES("entire_warehouse"="true")` gains a
  `RemoteSpillWriteSize` column (total row) as the billing input.

Tests: `SpillFileS3Test` (in-memory object store, 16 cases incl. budget,
cancellation and upload-failure paths), `MetaServiceTest.SpillStatsTest`,
`KeysTest.StatsSpillKeyTest`, recycler tests, and a column check in
`test_show_data_warehouse.groovy`.

Claude-Session: https://claude.ai/code/session_01Jrdwwwh8bZSnVwoCwykHse
…w of the S3 spill path

- Object root and spill stats record are keyed by the FE-assigned backend_id
  instead of cloud_unique_id, which every BE added by one ADD BACKEND
  statement shares. SpillStatsPB gains backend_id; a live BE whose id
  changes keeps the bound id and logs a warning.
- Startup cleanup discovers old boot generations through a per-boot marker
  object (spill/{backend_id}/boots/{boot_id}, refreshed daily) and deletes
  one generation per GC round instead of listing the whole BE prefix.
- Meta-service: report_spill_stats declares its get side for the KV bvars,
  rejects reports of an older boot_id (a late duplicate or a clock that went
  backwards) and never rolls back the totals of the current boot.
- FE getSpillStats goes through executeWithMetrics for retries/reconnect.
- Validators for spill_s3_storage_limit_bytes (>= 0) and
  spill_s3_max_inflight_upload_bytes (> 0); a budget below two upload
  buffers is a warning (degraded mode) at startup as at runtime; the recycler
  skips a non-positive spill_objects_expire_time_second.
- The upload budget over-release invariant is checked in release builds.
- Reporting documented as best-effort; multipart residue guidance added to
  the config comments.

Claude-Session: https://claude.ai/code/session_01Jrdwwwh8bZSnVwoCwykHse
Pure refactor, no behaviour change. The spill store was one class with an
`_is_remote` flag and a branch in every method. It is now:

- spill_data_dir.{h,cpp}: `SpillDataDir`, the abstract base owning what
  both kinds share (path, spill root, byte accounting against the limit,
  metrics), and `LocalSpillDataDir` (disk probing, spill_gc directory,
  usage-based disk selection).
- remote_spill_data_dir.{h,cpp}: `RemoteSpillDataDir` (storage vault
  binding, backend_id/boot_id, data and boot-marker key layout). Cloud
  headers are only included by its .cpp.

`SpillFileManager` keeps typed views of the stores (`_local_stores` and
at most one `_remote_store`) instead of testing `is_remote()` in every
loop; writers and readers still use the virtual `is_remote()` only to pick
local or remote counters.

Claude-Session: https://claude.ai/code/session_01Jrdwwwh8bZSnVwoCwykHse
…elled queries, drop spill of deleted instances

- The upload budget is charged with the allocated capacity of every
  submitted buffer (s3_write_buffer_size), not with its payload: a partially
  filled last buffer keeps its full allocation, so with parts smaller than
  the buffer the budget did not bound memory. The done callback reports the
  same capacity. A part size below the buffer size is warned at startup.
- SpillRemoteUploadBudget::acquire checks cancellation before admitting on
  the fast path and after every wake-up, so a cancelled query never starts
  a new upload.
- recycle_deleted_instance_data deletes the spill/ prefix on the
  snapshot-enabled path too, where only referenced rowsets are recycled
  selectively; spill objects are not referenced by anything.
- Tests: PartialBufferIsChargedAtCapacity, UploadBudgetRejectsCancelledQuery
  Immediately, spill objects in recycle_deleted_instance_with_orphan_tmp_rowset.

Claude-Session: https://claude.ai/code/session_01Jrdwwwh8bZSnVwoCwykHse
…ate with every buffer

- Spill objects live under spill/{instance_id}/{backend_id}/... . A storage
  vault can be shared by several instances (snapshot clones, rollback heirs)
  and backend ids are allocated per FE cluster, so the recycler of one
  instance must never touch another instance's spill. RemoteSpillDataDir
  learns the instance id through CloudMetaMgr::get_instance_id() when it
  binds; both recycler paths (TTL sweep and deleted-instance cleanup) act on
  spill/{instance_id}/ only, and the deleted-instance cleanup is restricted
  to S3/MOCK accessors like the TTL sweep.
- S3FileWriter::close() builds the buffer of an empty object itself so that
  it passes the upload gate like every other buffer; the done callback is
  now paired with the gate on every path.
- TQueryStatistics gains spill_write_bytes_to_remote_storage /
  spill_read_bytes_from_remote_storage, filled by the BE.
- Meta-service clears server-owned SpillStatsPB fields of a report before
  storing it; comments on the budget unit, cancellation and stale reports
  updated.
- Tests: GateAndDoneCallbackReportTheSameCapacity (partial PutObject,
  multipart with partial last buffer, empty object),
  CancelledQueryCloseIssuesNoUpload, other-instance objects in the recycler
  tests and in the startup-cleanup test.

Claude-Session: https://claude.ai/code/session_01Jrdwwwh8bZSnVwoCwykHse
…tead of the upload total

SHOW DATA reports current sizes, so the spill column must too. Each BE now
reports the spill bytes it holds in object storage right now
(RemoteSpillDataDir::get_spill_data_bytes()) instead of the bytes it has
uploaded since boot:

- SpillStatsPB carries remote_spill_bytes; the fold across boots
  (prior_boots_*) and the request count are gone. Meta-service replaces the
  record of a backend on every report of the same or a newer boot and still
  rejects an older boot. get_spill_stats sums the records refreshed within
  spill_objects_expire_time_second: an older record belongs to a BE that
  never came back, and the recycler has deleted its objects by then.
- The BE reports when the value changed (about once a minute), at least once
  an hour otherwise (heartbeat, so a live BE with stable spill is never
  mistaken for a dead one), always on its first report (replacing whatever
  the previous process of the same backend_id left behind) and once more at
  shutdown. Upload traffic stays available as BE metrics and bvars.
- SHOW DATA column renamed to RemoteSpillSize.
- MetaServiceTest.SpillStatsTest rewritten for the replace/expire semantics.

Claude-Session: https://claude.ai/code/session_01Jrdwwwh8bZSnVwoCwykHse
… fallback, bound the instance lookup

Follow-up to the local review of the current-size reporting:

- SpillStatsPB.report_seq orders the reports of one boot. A report whose
  client-side attempt timed out still executes on meta-service later; with
  the current-size semantics it would have overwritten a newer value (and
  the final 0 of a shutdown, which nobody corrects). Meta-service now rejects
  a report of the same boot with a smaller report_seq; a retried attempt
  carries the same seq and stays idempotent.
- S3FileWriter::_close_impl() no longer rebuilds the buffer of an empty
  object: close() builds it before passing the gate, so the fallback only
  ran for a writer whose gate refusal or failed append left no buffer, and
  it created the object and reported a done callback that never passed
  the gate. A failed empty writer now fails without creating the object.
- CloudMetaMgr::get_instance_id() is bounded to 2 attempts through the new
  _get_instance() shared with get_snapshot_properties(), and
  RemoteSpillDataDir::ensure_ready() does not repeat a failed lookup within
  meta_service_brpc_timeout_ms, so concurrent spill callers fail fast
  instead of each re-running the RPC under _init_mutex.
- The report heartbeat is wall-clock based (one hour) instead of counting GC
  rounds; the cadence check moved to _remote_gc; the config comment of
  spill_objects_expire_time_second states that it also expires reports.
- Comments updated to the current layout and semantics; dead
  clear_update_time_ms() removed.
- Tests: EmptyObjectRefusedByGateIsNotCreated, RemoteSpillStatsReportDecisions
  (first report, skip unchanged, failed report retried with a new seq,
  heartbeat), CancelledQueryCloseAbortsOpenMultipart; SpillStatsTest covers
  report_seq ordering and the shutdown case.

Claude-Session: https://claude.ai/code/session_01Jrdwwwh8bZSnVwoCwykHse
…rden the reporter test

- RemoteSpillDataDir::ensure_ready(): after a failed GetInstance the retry
  window is measured from the end of the RPC and lasts a minute (the GC
  thread's probe cadence). The previous window was sampled before the RPC
  and shorter than the RPC itself, so under a meta-service timeout it had
  expired by the time the next waiter took _init_mutex and the RPC ran
  back to back.
- RemoteSpillStatsReportDecisions: the recorder is owned by the hook
  (shared_ptr, mutex) and detached by a Defer, so an assertion failure or
  the final report of stop() cannot reach freed stack state; the GC thread
  is kept out of the test by a large spill_gc_interval_ms.
- FailedAppendThenCloseCreatesNoObject: a writer whose first append failed
  before any byte was counted creates no object when closed.
- The precondition of S3FileWriter::_close_impl() is a DORIS_CHECK with a
  message; report_seq and heartbeat comments made precise.

Claude-Session: https://claude.ai/code/session_01Jrdwwwh8bZSnVwoCwykHse
…he spill size for SHOW DATA

Spilling already runs on an IO thread that blocks on every append, and
S3FileWriter uploads buffers asynchronously within a part, so the only
thing the non-blocking close bought was the tail of the uploads plus one
CompleteMultipartUpload per part boundary: negligible against the
default 1 GB part. Drop the ClosingPart queue, _reap_closing_parts and
_finish_part; _close_current_part closes the writer synchronously and
reconciles the budget, the request statistics, the multipart abort and
the SpillFile registration in place. The ledger, statistics and abort
stay because they come from the asynchronous uploads inside a part, not
from the close.

Two S3 spill tests are adapted: UploadBudgetWaitIsCancellable used a
part exactly as large as the budget, so with a synchronous close the
writer waited for the blocked uploads instead of in the gate; it now
uses a larger part. CancelledQueryCloseAbortsOpenMultipart read its
PutObject baseline while uploads were still queued, a pre-existing race
with the upload threads; it now waits for the in-flight bytes to drain
first, so only the pending partial buffer is left for the gate to
refuse.

SHOW DATA's RemoteSpillSize was the only synchronous meta-service call
in that command. CloudTabletStatMgr, which every FE runs once per
tablet_stat_update_interval_second, now fetches get_spill_stats in its
own try block and SHOW DATA reads the cached value like the tablet
sizes. Because the column is a billing input, a value that was never
fetched or is older than the new fe.conf option
cloud_spill_stats_max_age_second (default 300 s) is reported as an
error instead of being shown as stale or zero.

Claude-Session: https://claude.ai/code/session_01CA1ot1uEwvpQzCJUhmdH32
_remote_gc, _report_remote_spill_stats, _remote_write_boot_marker and
_remote_startup_cleanup took the remote store as a parameter although
the manager owns exactly one; use the member directly and make
flush_remote_spill_stats an inline call of the final report. The header
no longer includes storage/options.h, io/fs/file_system.h and
<string_view>; the two tests that relied on them transitively include
what they use.

Claude-Session: https://claude.ai/code/session_01CA1ot1uEwvpQzCJUhmdH32
…drop the meta-service spill stats

The RemoteSpillSize column only has to reflect what the alive BEs hold:
objects left behind by a BE that is gone are removed by its restart or by
the recycler's TTL sweep and need not be counted meanwhile. With that
scope there is nothing a BE has to leave behind in meta-service, so the
whole report path goes: the periodic and shutdown reports of
SpillFileManager, CloudMetaMgr::report_spill_stats and its rate limiter,
the report_spill_stats/get_spill_stats RPCs, SpillStatsPB and its key,
bvars, the recycler's cleanup of that key range, the FE meta-service
client methods, and their tests.

Instead, get_be_resource (the brpc the FE already uses for admission
control) returns the spill bytes the BE currently holds in object
storage, and CloudTabletStatMgr polls every alive BE of every cluster
once per cycle and caches the sum for SHOW DATA. A poll with any failed
BE keeps the previous value rather than publishing a partial sum; the
existing cloud_spill_stats_max_age_second guard turns a poll that keeps
failing into an error instead of a stale number.

Claude-Session: https://claude.ai/code/session_01CA1ot1uEwvpQzCJUhmdH32
…the spill size as a SHOW DATA row

The audit log (fe.audit.log and __internal_schema.audit_log) now carries
SpillWriteBytesToRemoteStorage / SpillReadBytesFromRemoteStorage next to
the local spill bytes. The BE already reports them in TQueryStatistics;
WorkloadRuntimeStatusMgr sums them across BEs and fills them into the
audit event like the other byte counters. The audit table gets the two
columns after the local spill columns; an existing table picks them up
through the schema check at startup.

SHOW DATA PROPERTIES("entire_warehouse"="true") goes back to its original
four columns. In cloud mode the spill bytes currently held in object
storage are listed as a row of their own, __remote_spill__, and counted
in the total, so a consumer summing DataSize sees them without parsing a
new column. A user database name has to start with a letter, so the row
name cannot collide with one. The row is only added for the whole
warehouse, not when db_names restricts the listing.

Claude-Session: https://claude.ai/code/session_01CA1ot1uEwvpQzCJUhmdH32
…_id}/...

Remote spill objects were keyed by instance id, backend id and boot id
(spill/{instance_id}/{backend_id}/data/{boot_id}/{query_id}/... plus one
boot marker per process). They now live under spill/{host}/{query_id}/...,
where host is the address the BE advertises, so a key tells which BE and
which query wrote it. The store no longer waits for the FE heartbeat or
calls GetInstance to become ready; the GetInstance helper and the retry
cap added for it are reverted.

Without boot generations, the startup cleanup works on query
directories: a query registers its directory with SpillFileManager
before the first object of a spill file is written and unregisters it
when the query deletes its spill directory. Once the store is ready the
GC thread lists spill/{host}/ once, takes the directories nobody
registered as the residue of the previous process, and deletes one per
round. Directories that appear after that listing are left alone.

The layout assumes that no two BEs writing to the vault share an
address. The keys no longer name the instance, so the recycler sweeps
the whole spill/ prefix by age, and the snapshot-enabled deleted-instance
path runs that age-based sweep instead of dropping a prefix.

Claude-Session: https://claude.ai/code/session_01CA1ot1uEwvpQzCJUhmdH32
With spill/{ip}/ two BEs on one host shared a directory: the same query
has the same directory on both, so the BE that finished first deleted
the other's spill, and a restarting BE took the other's live directories
for residue. The directory is now spill/{ip}_{port}/, port being the
heartbeat_service_port, the pair FE identifies a BE by.

Claude-Session: https://claude.ai/code/session_01CA1ot1uEwvpQzCJUhmdH32
…tate is destroyed

~RuntimeState cleared _obj_pool in its body, before its members were
destroyed. Operator local states are members, and they hold counters of
profiles allocated in that pool, so a local state whose destructor still
used them read freed memory. A spill writer that was left open closes its
last part in its destructor and updates its write timer: cancelling a
query that spills to object storage makes the writer's close fail, the
hash join sink stops closing the remaining writers at the first error,
and the BE crashed tearing down the task.

Release the local states first, in the order the members would be
destroyed, and then the pool.

Claude-Session: https://claude.ai/code/session_01CA1ot1uEwvpQzCJUhmdH32
…ng remote spill files

SpillFileReader issued one read per spilled block plus three reads per part
footer. On object storage every read is a GET, so reading spill back was
bound by request latency: a TPC-DS q11 run at sf10000 read 572 GB of spill
with 1.59M GETs of ~377 KB each (~14 MB/s per stream), while the write path
uploaded the same bytes in 5 MB parts.

- Serve blocks from a window: on a miss, read the block and the following
  adjacent blocks in one read of at most spill_s3_read_coalesce_bytes
  (default 8 MB, capped by spill_buffer_size_bytes). A block larger than the
  window is still read whole. 0 restores exact block-by-block reads.
- Read the part footer with one tail read of 64 KB, and fetch a part that
  fits in one window whole when it is opened for reading.
- Seeking within the current part reuses the buffered bytes; parts skipped by
  a seek only read their footers.
- Return an error on short reads and corrupted block counts or offsets
  instead of relying on DCHECKs.

Local disk keeps reading every block separately; its footer now takes two
exact reads instead of three.

Claude-Session: https://claude.ai/code/session_01EFY2NVjqRqSektbK1HpTc8
The spill size shown by SHOW DATA was refreshed at the start of every
CloudTabletStatMgr round, so its freshness depended on how long a round of
the tablet stats took. With tablet_stat_update_interval_second at 600, or a
slow round on a large instance, the value aged past
cloud_spill_stats_max_age_second (300) every round and the whole
"SHOW DATA ... entire_warehouse" command failed, also on clusters that never
spill to object storage.

- Move the polling into RemoteSpillStatsPoller, a daemon started on every FE
  with its own interval, cloud_spill_stats_poll_interval_second (60 s,
  mutable). CloudTabletStatMgr is back to its original form.
- Refuse a value only when it is older than the larger of
  cloud_spill_stats_max_age_second and three poll intervals, so a long
  interval cannot make every value stale.

Claude-Session: https://claude.ai/code/session_01EFY2NVjqRqSektbK1HpTc8
…deleted spill charged

Recycler: the sweep deleted every spill object older than
spill_objects_expire_time_second, so a legal query running longer than the
TTL (7 days by default; query_timeout has no upper bound) could lose data it
still needed.

- A BE with a remote spill store rewrites spill/{ip}_{port}/_heartbeat every
  spill_s3_heartbeat_interval_second (1 hour by default).
- The recycler groups the objects under spill/ by BE directory and deletes a
  directory only when nothing in it, the heartbeat included, changed for the
  whole TTL. Objects of a live BE are kept however old; the expiration time
  still applies to the deletion so objects written after the listing survive.
- MockAccessor now records modification times and honours the expiration
  time of delete_prefix, so the test covers live, dead and active BEs.

BE: when deleting a spill file failed, its bytes were released from
spill_s3_storage_limit_bytes although the objects were still stored, and the
pending query-directory retry kept only a path. A long deletion outage let
every failed query free and reuse the capacity while the retained objects
grew without bound, and the spill size reported to SHOW DATA dropped.

- SpillFile::gc() hands a failed deletion with its bytes to the manager, which
  keeps them charged until a retry deletes the objects.
- Pending deletions are merged by directory: a query directory absorbs the
  pending spill file directories under it, so an outage keeps about one entry
  per query.

Claude-Session: https://claude.ai/code/session_01EFY2NVjqRqSektbK1HpTc8
…he spill sweep in the recycler statistics

Recycler: the sweep grouped the listed objects by "spill/{ip}_{port}/" and
deleted every group whose newest object was older than the TTL. An object
whose key is exactly "spill/" (the directory marker the S3 console creates
for a "folder") became a group of its own and, because the group ended with
'/', was deleted with delete_prefix("spill/", expiration): that removed every
object older than the TTL of every BE directory, the parts of a running long
query included, which is exactly what the heartbeat protects against.

- A group is deleted as a prefix only when it is a BE directory; an object
  directly under "spill/" is deleted alone.
- The listing is shared with a new scan_and_statistics_expired_spill_objects,
  so the sweep shows up in the recycler statistics like its siblings, and the
  sweep publishes the need-recycle numbers from its own listing.
- MockAccessor lists object sizes; the recycler test adds the "spill/"
  marker and a fresh stray object and checks the expired groups with their
  bytes.

BE: _retry_pending_spill_directories (renamed, it retries spill file
directories too), the retry path and SpillFile::gc() assert instead of
branching on a store that cannot be null or not ready (the store is ready
before any spill file exists); the heartbeat config documents that
interval 0 leaves a BE directory protected only by its newest object. Tests add the heartbeat
failure retry and the query directory deleted while a pending spill file
directory under it is still queued.

FE: the audit schema row test is folded into the one master already has.

Claude-Session: https://claude.ai/code/session_01MQABk3ZCiHWQYzwLr7NL2e
@mrhhsg
mrhhsg force-pushed the feat/cloud-spill-s3 branch from 94124f3 to 71a21e2 Compare September 24, 2026 07:55
@mrhhsg

mrhhsg commented Sep 24, 2026

Copy link
Copy Markdown
Member Author

/review-light

@github-actions github-actions Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Request changes. Static review of exact head 71a21e2af8850b8467c997d9f0adef13cdf8d010 found nine distinct issues: five P1 correctness/scalability issues and four P2 observability/freshness gaps. Existing inline threads were treated as hard duplicate fences.

Critical-checkpoint conclusions:

  • Goal and tests: the PR adds cloud S3 spill storage, lifecycle cleanup, accounting/billing, and remote-spill observability. It has substantial BE, cloud, FE, and regression coverage, but the goal is not safely complete because shared-vault ownership, deleted-instance cleanup, terminal write-error propagation, error-unwind accounting, and recycler scalability are not covered.
  • Scope: the 62-file change is broad but cohesive. The new storage path itself is understandable; parallel consumer and producer paths were missed (MAIN-3, MAIN-5, MAIN-7, MAIN-9).
  • Concurrency and lifecycle: I traced GC, recycler, scheduled FE polling, async S3 upload callbacks, upload-budget admission/release, and RuntimeState/QueryContext teardown. The ordinary gate/callback branches and lock ordering appear balanced with no new deadlock found, but MAIN-1, MAIN-2, MAIN-4, and MAIN-9 expose cross-instance, teardown, or terminal-close lifecycle failures.
  • Configuration and compatibility: mutable limits, upload budget, heartbeat/TTL, poll interval, and max age are read on live paths; startup/dynamic parity already has an existing review thread. Appended optional Thrift/protobuf fields and name-mapped audit columns preserve rolling compatibility. MAIN-8 violates the configured freshness guarantee after a wall-clock rollback.
  • Conditions and parallel paths: TTL filtering is the direct trigger for MAIN-2. The local-spill analogues exposed missing live-stat aggregation, Iceberg counter registration, HTTP-plan audit initialization, and Iceberg terminal-close handling.
  • Tests/results: added tests cover core S3 writer/reader/quota behavior, cloud expiry, FE polling/audit aggregation, and SHOW DATA, and the checked-in audit schema output is consistent. Missing negative coverage is called out inline: shared instances/endpoints, default-TTL instance deletion, repartition error unwinding, multi-page shared-vault scans, Iceberg S3 close failures/counters, live consumers, HTTP-plan audit output, and clock rollback. No build or test was run for this review, as required by the task.
  • Observability: the PR adds useful counters, metrics, audit fields, and logs, but MAIN-3, MAIN-5, and MAIN-7 leave supported paths invisible or report sentinel values.
  • Persistence and data writes: no new FE EditLog path is involved. Remote multipart writes and cloud recycler state were traced through failure/cleanup boundaries; MAIN-2 loses the durable cleanup owner, MAIN-4 loses the accounting owner, and MAIN-9 suppresses a terminal write failure after intermediate inputs may already be deleted. Crash-before-report, failed delete, multipart residue, retry-attempt counting, and heartbeat deletion races are already covered by existing threads and were not duplicated.
  • FE/BE contracts and performance: the new optional wire fields are compatible, but live consumers are incomplete (MAIN-3). MAIN-6 adds repeated unbounded N * M object listings on shared vaults. The fixed 64-KiB footer probe is bounded and was not considered a separate high-value issue.
  • Other checkpoints: no new nullable-column, visible-version, MoW, static-initialization, or schema-format issue applies here. Error handling and ownership were rechecked; MAIN-4 and MAIN-9 are the remaining ownership/error-propagation defects.

The focus file contains only -light and names no subsystem; I kept the presentation concise but completed the full review. The final BE pass found MAIN-9 at the three-round cap, so this review is capped/incomplete rather than converged. Validation is static only.

Comment thread be/src/exec/spill/remote_spill_data_dir.cpp
Comment thread cloud/src/recycler/recycler.cpp Outdated
Comment thread be/src/runtime/workload_management/resource_context.cpp
Comment thread be/src/exec/spill/spill_file_writer.cpp
Comment thread be/src/exec/operator/spill_counters.h
Comment thread cloud/src/recycler/recycler.cpp Outdated
Comment thread fe/fe-core/src/main/java/org/apache/doris/plugin/AuditEvent.java
Comment thread be/src/exec/spill/spill_file_writer.cpp
### What problem does this PR solve?

Issue Number: None

Related PR: apache#68032

Problem Summary: When reading a remote spill part, the footer is fetched by one probe of the
part's tail. The probe was always min(part size, 64 KiB), even when
spill_s3_read_coalesce_bytes (or the query's spill buffer size) was smaller. So a 1 KiB setting
still issued a 40-64 KiB GET and allocated a buffer of that size, which breaks the documented
per-GET bound. The probe is now bounded by the effective coalesced read size as well, but it is
never smaller than the 16-byte tail. If the offsets array does not fit, it is fetched separately,
as before.

Also document that an empty spill_s3_storage_vault is resolved to the default vault when the BE
spills for the first time. A later SET DEFAULT STORAGE VAULT applies to spill after the BE
restarts. Also document that the rest of an offset array that does not fit in the probe is read
whole, in the same way as a single block that is larger than the read size.

### Release note

None

### Check List (For Author)

- Test: Unit Test
    - SpillFileS3Test.ReadCoalescesAdjacentBlocks now expects the probe to be capped by a 1 KiB
      and a 40 KiB window; new FooterLargerThanProbe covers a window smaller than the block
      offset array (probe, then one GET for the rest of the footer); SpillFile* (63 tests) passed
- Behavior changed: No
- Does this need documentation: No

Claude-Session: https://claude.ai/code/session_01B5f3dfnSMXhasgsShyagfR
@mrhhsg

mrhhsg commented Sep 30, 2026

Copy link
Copy Markdown
Member Author

run buildall

@hello-stephen

Copy link
Copy Markdown
Contributor

Cloud UT Coverage Report

Increment line coverage 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 77.54% (2064/2662)
Line Coverage 65.91% (38015/57678)
Region Coverage 53.43% (35695/66810)
Branch Coverage 56.79% (11464/20186)

@hello-stephen

Copy link
Copy Markdown
Contributor
TPC-H: Total hot run time: 28082 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpch-tools
Tpch sf100 test result on commit fcc529eadf1f84554ef3d2f11d4da67a3f691a08, data reload: false

------ Round 1 ----------------------------------
============================================
q1	17401	3996	4047	3996
q2	2180	361	310	310
q3	10016	1367	789	789
q4	4686	477	355	355
q5	7508	833	551	551
q6	184	187	141	141
q7	737	800	614	614
q8	9354	1530	1590	1530
q9	5830	4230	4190	4190
q10	6879	1314	1035	1035
q11	444	270	251	251
q12	635	415	292	292
q13	18103	2628	2008	2008
q14	272	261	238	238
q15	q16	765	733	663	663
q17	1652	1161	1000	1000
q18	6531	5623	5528	5528
q19	1281	1305	1098	1098
q20	497	405	263	263
q21	5410	3427	2917	2917
q22	437	376	313	313
Total cold run time: 100802 ms
Total hot run time: 28082 ms

----- Round 2, with runtime_filter_mode=off -----
============================================
q1	4715	4536	4637	4536
q2	705	564	565	564
q3	4781	5354	4718	4718
q4	2215	2357	1453	1453
q5	4734	4457	4520	4457
q6	236	174	127	127
q7	1851	1706	1517	1517
q8	2352	2138	2139	2138
q9	7247	6913	6856	6856
q10	3615	3554	3088	3088
q11	515	399	354	354
q12	699	698	507	507
q13	2299	2604	2026	2026
q14	267	280	249	249
q15	q16	658	691	600	600
q17	7323	6803	6658	6658
q18	11891	11051	11813	11051
q19	1094	959	1023	959
q20	2212	2195	1932	1932
q21	5063	4228	4424	4228
q22	510	454	400	400
Total cold run time: 64982 ms
Total hot run time: 58418 ms

@hello-stephen

Copy link
Copy Markdown
Contributor
TPC-DS: Total hot run time: 152654 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpcds-tools
TPC-DS sf100 test result on commit fcc529eadf1f84554ef3d2f11d4da67a3f691a08, data reload: false

query5	4317	591	461	461
query6	433	214	191	191
query7	4819	533	297	297
query8	333	180	181	180
query9	8801	3964	3974	3964
query10	450	316	259	259
query11	5937	3535	3221	3221
query12	143	87	85	85
query13	1269	605	407	407
query14	6529	4530	4246	4246
query14_1	3992	3994	4006	3994
query15	202	198	181	181
query16	989	458	412	412
query17	928	666	535	535
query18	2433	463	336	336
query19	210	184	162	162
query20	84	79	84	79
query21	226	134	116	116
query22	13023	13030	12833	12833
query23	13910	12906	12434	12434
query23_1	12614	12491	12486	12486
query24	7222	1161	660	660
query24_1	695	750	706	706
query25	569	443	371	371
query26	1278	318	164	164
query27	2698	526	337	337
query28	4579	1946	1940	1940
query29	1630	730	541	541
query30	297	225	185	185
query31	888	776	635	635
query32	152	100	97	97
query33	563	314	255	255
query34	1169	1102	605	605
query35	738	745	642	642
query36	805	797	698	698
query37	150	109	96	96
query38	1827	1757	1701	1701
query39	687	666	673	666
query39_1	653	643	659	643
query40	224	127	101	101
query41	79	73	69	69
query42	100	96	92	92
query43	350	353	306	306
query44	1403	714	711	711
query45	188	174	162	162
query46	1084	1183	714	714
query47	1475	1513	1397	1397
query48	392	392	356	356
query49	579	394	294	294
query50	987	360	266	266
query51	10519	10624	10290	10290
query52	89	96	76	76
query53	236	254	175	175
query54	249	207	201	201
query55	78	76	74	74
query56	232	199	216	199
query57	1542	1416	1304	1304
query58	279	259	257	257
query59	1983	2049	1879	1879
query60	299	241	228	228
query61	145	160	145	145
query62	404	319	266	266
query63	219	180	184	180
query64	2821	1022	819	819
query65	3519	3392	3431	3392
query66	1790	410	297	297
query67	20013	20149	20189	20149
query68	3173	1543	870	870
query69	403	292	257	257
query70	900	807	806	806
query71	290	236	217	217
query72	2646	2571	2239	2239
query73	819	788	423	423
query74	4609	4481	4299	4299
query75	2285	2303	1931	1931
query76	2311	1145	743	743
query77	365	388	310	310
query78	9108	9157	8462	8462
query79	1363	1222	751	751
query80	580	464	361	361
query81	541	325	282	282
query82	660	166	131	131
query83	320	225	197	197
query84	329	147	117	117
query85	850	458	380	380
query86	332	248	235	235
query87	2005	1941	1833	1833
query88	3679	2728	2668	2668
query89	374	305	249	249
query90	1907	187	181	181
query91	172	157	127	127
query92	106	89	86	86
query93	1444	1435	854	854
query94	543	353	301	301
query95	664	373	445	373
query96	1023	821	325	325
query97	2441	2395	2322	2322
query98	163	152	145	145
query99	712	732	611	611
Total cold run time: 236353 ms
Total hot run time: 152654 ms

@hello-stephen

Copy link
Copy Markdown
Contributor
ClickBench: Total hot run time: 25.23 s
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/clickbench-tools
ClickBench test result on commit fcc529eadf1f84554ef3d2f11d4da67a3f691a08, data reload: false

query1	0.01	0.01	0.00
query2	0.13	0.09	0.08
query3	0.38	0.24	0.24
query4	1.62	0.25	0.25
query5	0.34	0.32	0.32
query6	1.16	0.66	0.68
query7	0.04	0.01	0.00
query8	0.10	0.07	0.07
query9	0.50	0.41	0.41
query10	0.62	0.58	0.59
query11	0.33	0.19	0.20
query12	0.33	0.19	0.19
query13	0.54	0.52	0.52
query14	0.87	0.86	0.87
query15	0.67	0.60	0.60
query16	0.39	0.39	0.38
query17	1.01	1.02	1.01
query18	0.30	0.29	0.29
query19	1.92	1.82	1.81
query20	0.02	0.01	0.02
query21	15.51	0.38	0.30
query22	4.86	0.12	0.12
query23	15.85	0.48	0.29
query24	2.36	0.61	0.42
query25	0.16	0.10	0.11
query26	0.75	0.26	0.21
query27	0.10	0.10	0.10
query28	3.48	0.89	0.45
query29	12.51	4.26	3.30
query30	0.38	0.27	0.25
query31	2.77	0.58	0.34
query32	3.23	0.60	0.48
query33	2.93	2.91	3.01
query34	15.81	3.94	3.28
query35	3.19	3.16	3.17
query36	0.65	0.54	0.51
query37	0.12	0.09	0.10
query38	0.08	0.06	0.07
query39	0.06	0.07	0.06
query40	0.20	0.16	0.16
query41	0.10	0.05	0.06
query42	0.06	0.06	0.06
query43	0.07	0.06	0.06
Total cold run time: 96.51 s
Total hot run time: 25.23 s

@hello-stephen

Copy link
Copy Markdown
Contributor

BE Regression && UT Coverage Report

Increment line coverage 69.90% (72/103) 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 76.57% (34964/45662)
Line Coverage 61.73% (394604/639248)
Region Coverage 58.12% (332285/571739)
Branch Coverage 59.00% (152439/258381)

### What problem does this PR solve?

Issue Number: None

Related PR: apache#68032

Problem Summary: Remote spill lacked deterministic coverage for the FE polling RPC's live-backend sum and failure retention, the nonzero SHOW DATA warehouse spill row and total, and the BE-to-FE query-statistics transfer. Add focused tests for successful and failed polls, stale billing values, cloud/non-cloud and database-restricted listings, remote store identity, and local-versus-remote spill counters.

### Release note

None

### Check List (For Author)

- Test: Unit Test: FE remote spill, SHOW DATA, audit, current-query and workload statistics (50 tests); BE SpillFileS3Test (28 tests). FE build, C++ format and build hygiene checks passed. Clang-tidy attempted but blocked by toolchain stddef.h and pre-existing header diagnostics.
- Behavior changed: No. Test-only change.
- Does this need documentation: No.
@mrhhsg

mrhhsg commented Oct 7, 2026

Copy link
Copy Markdown
Member Author

run buildall

@mrhhsg

mrhhsg commented Oct 7, 2026

Copy link
Copy Markdown
Member Author

/review-light

@github-actions github-actions Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Complete static review of PR #68032 at 7ad11da4c1852b91b141b66022e97fbfb7958108. Changes requested: the already reported P1 default-vault rotation issue still applies. There are no new inline findings to repost.

The following existing P2 threads also still apply: reserved spill bytes reported as stored bytes, a reachable BE skipped after heartbeat failure, remote audit bytes lost when the final BE report is late, and mixed-version FE current-query rows with different widths. The earlier footer-probe concern is fixed for its reported case at this head; older recycler and writer-close P1 threads refer to removed or repaired paths.

Critical checkpoints

  • Goal, scope, and focus: The change adds cloud S3 spill, bounded uploads, coalesced reads, per-query remote counters, and a live-BE SHOW DATA total. The 64 changed paths and their tests were reviewed. The supplied focus was -light; it yielded no separate issue. The implementation is large because it crosses BE I/O, cloud vaults, FE reporting, and protocol fields, but the changes follow that feature path.
  • Concurrency: S3 upload workers, the GC thread, query teardown, and FE polling are the relevant concurrent paths. Budget accounting uses a mutex; upload callbacks signal completion before final writer status; FE publishes an immutable volatile poll sample. No additional lock-order, callback lifetime, or deadlock issue was substantiated.
  • Lifecycle and static initialization: Spill-file GC, failed-delete retry, multipart abort, manager shutdown, and RuntimeState local-state teardown were traced. No new ownership cycle or cross-translation-unit initialization dependency was found. The existing P1 captures the default vault binding lifecycle error.
  • Configuration and conditions: BE storage type and explicit vault are startup settings; mutable capacity, upload-budget, and read-coalescing limits are observed by their respective paths. FE poll interval and age settings are read on later cycles. The default-vault setting does not follow a live default rotation, as raised in the P1. Error/cancellation paths were checked; no other distinct condition or silent failure was substantiated.
  • Compatibility and parallel paths: The new protobuf and Thrift fields are optional and old BEs cannot produce this new S3 spill. Iceberg spill closes its writer and registers remote counters; the HTTP-plan audit path initializes the new fields. The existing P2 mixed-FE row-shape and delayed-audit-report gaps remain. No other parallel FE/BE transport omission was found.
  • Tests and results: Added BE roundtrip, capacity, upload-failure, cancellation, read-coalescing, GC, and lifecycle tests, FE polling/SHOW DATA/audit/stat tests, and regression changes were inspected. The new audit expected-schema rows match the added columns on static inspection. Heartbeat-failure, delayed-final-report, mixed-FE, and vault-rotation cases remain the already reported gaps. Builds, tests, and final ELF ABI checks were not run under this static-only review contract; inspected tests are not independent validation.
  • Observability, persistence, and writes: Remote spill counters, bvars, metrics, audit columns, and get_be_resource reporting are wired. The two existing P2 live-total discrepancies remain. No new Doris transaction or EditLog state is introduced. S3 part publication and cleanup were traced through failures; crash residue is explicitly left to the bucket lifecycle rule and excluded from the documented live-BE total.
  • Performance and remaining issues: Upload admission and GET coalescing address the new hot paths. No additional substantiated performance or data-correctness issue emerged from two complete review passes and the separate risk-focused pass. All proposed candidates were verified and either matched an existing thread or were dismissed with code evidence; no candidate remains unresolved.

Existing P0/P1 findings confirmed for this head: #68032 (comment)

@hello-stephen

Copy link
Copy Markdown
Contributor

Cloud UT Coverage Report

Increment line coverage 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 77.79% (2073/2665)
Line Coverage 65.97% (38076/57716)
Region Coverage 53.44% (35715/66831)
Branch Coverage 56.96% (11508/20204)

@hello-stephen

Copy link
Copy Markdown
Contributor
TPC-H: Total hot run time: 29210 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpch-tools
Tpch sf100 test result on commit 7ad11da4c1852b91b141b66022e97fbfb7958108, data reload: false

------ Round 1 ----------------------------------
============================================
q1	17382	5636	5585	5585
q2	2100	323	256	256
q3	10137	1402	850	850
q4	4671	475	359	359
q5	7453	806	531	531
q6	189	183	158	158
q7	787	790	606	606
q8	9251	1357	1060	1060
q9	5702	4503	4451	4451
q10	6879	1300	1033	1033
q11	437	252	235	235
q12	629	402	286	286
q13	18093	3067	2382	2382
q14	277	278	245	245
q15	q16	743	723	682	682
q17	1206	740	590	590
q18	7304	6236	6218	6218
q19	1109	965	597	597
q20	387	337	223	223
q21	5834	2568	2619	2568
q22	397	403	295	295
Total cold run time: 100967 ms
Total hot run time: 29210 ms

----- Round 2, with runtime_filter_mode=off -----
============================================
q1	6547	6489	6514	6489
q2	750	609	567	567
q3	5540	5665	5320	5320
q4	2296	2166	1596	1596
q5	5859	5778	5801	5778
q6	241	184	141	141
q7	2280	2056	1814	1814
q8	3152	2752	2764	2752
q9	8275	7866	7744	7744
q10	3677	3596	3224	3224
q11	604	413	372	372
q12	659	702	502	502
q13	2742	3066	2378	2378
q14	278	293	279	279
q15	q16	671	677	635	635
q17	7852	7028	6944	6944
q18	13109	12228	13072	12228
q19	923	783	792	783
q20	2196	2148	1946	1946
q21	5763	4673	4837	4673
q22	511	454	415	415
Total cold run time: 73925 ms
Total hot run time: 66580 ms

@hello-stephen

Copy link
Copy Markdown
Contributor
TPC-DS: Total hot run time: 151779 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpcds-tools
TPC-DS sf100 test result on commit 7ad11da4c1852b91b141b66022e97fbfb7958108, data reload: false

query5	4303	615	464	464
query6	430	206	190	190
query7	4806	447	213	213
query8	317	181	162	162
query9	8780	4004	3956	3956
query10	469	319	254	254
query11	6037	3427	3198	3198
query12	141	90	89	89
query13	1256	423	307	307
query14	6544	4782	4526	4526
query14_1	4252	4242	4241	4241
query15	203	196	186	186
query16	947	477	465	465
query17	845	682	544	544
query18	2427	435	306	306
query19	191	160	135	135
query20	83	80	78	78
query21	216	135	115	115
query22	12992	12966	12841	12841
query23	13208	12609	12069	12069
query23_1	12365	12124	12364	12124
query24	7050	816	461	461
query24_1	477	455	458	455
query25	499	389	335	335
query26	1244	248	136	136
query27	2781	435	264	264
query28	4606	1915	1911	1911
query29	1515	573	433	433
query30	294	208	180	180
query31	885	759	639	639
query32	133	104	86	86
query33	502	294	237	237
query34	933	857	500	500
query35	729	767	639	639
query36	860	819	750	750
query37	125	98	95	95
query38	1777	1726	1683	1683
query39	715	697	683	683
query39_1	669	663	675	663
query40	216	125	94	94
query41	66	63	63	63
query42	90	85	82	82
query43	359	380	326	326
query44	1274	688	673	673
query45	183	171	166	166
query46	820	942	572	572
query47	2952	3013	2827	2827
query48	316	305	211	211
query49	589	397	296	296
query50	674	267	214	214
query51	10278	10210	10145	10145
query52	78	81	72	72
query53	195	203	150	150
query54	230	213	224	213
query55	73	68	66	66
query56	235	215	218	215
query57	1589	1577	1535	1535
query58	293	255	250	250
query59	2370	2367	2117	2117
query60	276	240	219	219
query61	142	149	147	147
query62	383	344	289	289
query63	197	158	165	158
query64	2733	943	801	801
query65	3441	3382	3405	3382
query66	1773	409	317	317
query67	20503	20167	20144	20144
query68	3021	970	603	603
query69	381	296	246	246
query70	912	848	835	835
query71	291	234	220	220
query72	2666	2564	2471	2471
query73	506	528	292	292
query74	4564	4480	4317	4317
query75	2336	2291	1936	1936
query76	2259	1026	608	608
query77	362	396	305	305
query78	9098	8977	8474	8474
query79	878	849	526	526
query80	466	425	356	356
query81	521	320	287	287
query82	209	136	115	115
query83	194	196	185	185
query84	283	120	101	101
query85	806	416	367	367
query86	256	240	228	228
query87	2035	1939	1824	1824
query88	3588	2703	2678	2678
query89	324	289	262	262
query90	2066	184	174	174
query91	154	153	129	129
query92	96	83	87	83
query93	976	974	569	569
query94	449	323	279	279
query95	580	395	334	334
query96	660	524	231	231
query97	2451	2429	2314	2314
query98	160	147	143	143
query99	773	773	656	656
Total cold run time: 232950 ms
Total hot run time: 151779 ms

@hello-stephen

Copy link
Copy Markdown
Contributor
ClickBench: Total hot run time: 25.44 s
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/clickbench-tools
ClickBench test result on commit 7ad11da4c1852b91b141b66022e97fbfb7958108, data reload: false

query1	0.00	0.01	0.01
query2	0.15	0.07	0.08
query3	0.36	0.21	0.21
query4	1.61	0.22	0.21
query5	0.31	0.28	0.28
query6	1.16	0.66	0.65
query7	0.04	0.01	0.01
query8	0.08	0.07	0.08
query9	0.48	0.39	0.39
query10	0.57	0.57	0.56
query11	0.31	0.19	0.18
query12	0.32	0.19	0.18
query13	0.51	0.52	0.51
query14	0.90	0.89	0.90
query15	0.67	0.58	0.58
query16	0.36	0.36	0.36
query17	1.01	1.01	1.01
query18	0.31	0.29	0.28
query19	1.91	1.86	1.94
query20	0.02	0.02	0.02
query21	15.45	0.36	0.31
query22	4.79	0.12	0.14
query23	15.84	0.49	0.29
query24	2.21	0.55	0.40
query25	0.15	0.11	0.10
query26	0.73	0.27	0.21
query27	0.09	0.10	0.10
query28	3.57	0.85	0.42
query29	12.45	4.26	3.31
query30	0.37	0.23	0.21
query31	2.77	0.59	0.34
query32	3.25	0.63	0.50
query33	2.99	3.05	2.95
query34	15.84	4.10	3.45
query35	3.39	3.39	3.38
query36	0.57	0.49	0.50
query37	0.13	0.10	0.09
query38	0.07	0.07	0.06
query39	0.06	0.06	0.05
query40	0.18	0.17	0.16
query41	0.11	0.05	0.05
query42	0.06	0.06	0.06
query43	0.06	0.06	0.06
Total cold run time: 96.21 s
Total hot run time: 25.44 s

@hello-stephen

Copy link
Copy Markdown
Contributor

FE UT Coverage Report

Increment line coverage 84.54% (82/97) 🎉
Increment coverage report
Complete coverage report

@hello-stephen

Copy link
Copy Markdown
Contributor

BE UT Coverage Report

Increment line coverage 88.00% (66/75) 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 65.11% (30848/47375)
Line Coverage 49.88% (323888/649281)
Region Coverage 45.37% (261822/577129)
Branch Coverage 47.02% (122740/261029)

@hello-stephen

Copy link
Copy Markdown
Contributor

FE Regression Coverage Report

Increment line coverage 77.32% (75/97) 🎉
Increment coverage report
Complete coverage report

@hello-stephen

Copy link
Copy Markdown
Contributor

BE Regression && UT Coverage Report

Increment line coverage 71.57% (73/102) 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 76.73% (35181/45848)
Line Coverage 61.93% (397503/641895)
Region Coverage 58.25% (334857/574853)
Branch Coverage 59.20% (154056/260240)

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants