Repository navigation
Conversation
|
Thank you for your contribution to Apache Doris. Please clearly describe your PR:
|
|
/review |
There was a problem hiding this comment.
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):
- 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.
- 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.
- 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).
- 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.
- 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).
- 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.
- 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).
- 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).
- 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.
- 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.
- 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.
- 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.
- 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).
- 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).
- 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).
- 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.
|
Thanks for the review. Pushed a follow-up commit addressing it; per finding: Fixed
Kept as designed (happy to discuss)
|
add6dfd to
02ba809
Compare
Local pipeline review — ✅ PASSschema: 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: trueNotes for maintainers
Reviewed locally with the |
0389435 to
94124f3
Compare
… 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
94124f3 to
71a21e2
Compare
|
/review-light |
There was a problem hiding this comment.
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 * Mobject 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.
### 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
|
run buildall |
Cloud UT Coverage ReportIncrement line coverage Increment coverage report
|
TPC-H: Total hot run time: 28082 ms |
TPC-DS: Total hot run time: 152654 ms |
ClickBench: Total hot run time: 25.23 s |
BE Regression && UT Coverage ReportIncrement line coverage Increment coverage report
|
### 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.
|
run buildall |
|
/review-light |
There was a problem hiding this comment.
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 DATAtotal. 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
RuntimeStatelocal-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_resourcereporting 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)
Cloud UT Coverage ReportIncrement line coverage Increment coverage report
|
TPC-H: Total hot run time: 29210 ms |
TPC-DS: Total hot run time: 151779 ms |
ClickBench: Total hot run time: 25.44 s |
FE UT Coverage ReportIncrement line coverage |
BE UT Coverage ReportIncrement line coverage Increment coverage report
|
FE Regression Coverage ReportIncrement line coverage |
BE Regression && UT Coverage ReportIncrement line coverage Increment coverage report
|
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 andSHOW DATAshows the sum.BE
RemoteSpillDataDirbound to the storage vault (spill_s3_storage_vault, default: the instance's default vault); objects live underspill/{ip}_{port}/{query_id}/{spill file}/{part}, whereipis the address the BE advertises andportitsheartbeat_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 listsspill/{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/SpillFileReaderreuseS3FileWriter/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;FileWriterOptionsgains an upload gate + done callback for this, honoured byS3FileWriter. The spill store is split intoSpillDataDir(base),LocalSpillDataDirandRemoteSpillDataDir.spill_remote_{read,write}_bytes,spill_remote_{get,put}_requestsand per-second throughput/QPS; profile countersSpillRemote*.TQueryStatisticsgainsspill_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.RuntimeStatereleases 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
recycle_expired_spill_objectsremoves objects underspill/older thanspill_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; anAbortIncompleteMultipartUploadbucket 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 thetotalrow. The row is added in cloud mode for the whole warehouse only, not whendb_namesrestricts the listing, and the column layout is unchanged. Every FE polls all alive BEs withget_be_resourceonce pertablet_stat_update_interval_second(CloudTabletStatMgr, like the tablet sizes) andSHOW DATAreads the cached sum; a poll with any failed BE keeps the previous value, and a value that is missing or older thancloud_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.fe.audit.logand__internal_schema.audit_log) gainsSpillWriteBytesToRemoteStorage/SpillReadBytesFromRemoteStorage(columnsspill_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 = s3in 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
Behavior changed:
spill_storage_type,spill_s3_storage_vault,spill_s3_storage_limit_bytes,spill_s3_max_inflight_upload_bytes), new meta-service configspill_objects_expire_time_second, new fe.conf optioncloud_spill_stats_max_age_second, new__remote_spill__row inSHOW DATA ... entire_warehouse(included in the total), two new audit log columns.Does this need documentation?
Check List (For Reviewer who merge this PR)
https://claude.ai/code/session_01Jrdwwwh8bZSnVwoCwykHse
https://claude.ai/code/session_01CA1ot1uEwvpQzCJUhmdH32