Repository navigation
Let each scatter shard decide the file cache against its own node (T-630) - #454
Merged
Merged
Conversation
…630)
T-629 made the coordinator decide for every shard of a scattered query,
from its own node's cache index. The shards read on other nodes, and
their blocks land in those nodes' caches, which that index never sees.
So a table warmed through its shards still looked cold to the
coordinator, and every auto run bypassed the cache again, on every
shard. A `fileCache: true` warming run could not change that.
Each shard now decides for itself (FileCache.shard_decision/4):
- The coordinator sends the job's mode (auto, true, false) as
file_cache_mode, and node_shards, how many of the job's shards run on
that node.
- A sealed shard unit carries its "bytes". On auto, the shard sums the
bytes of its sealed files that its own node's index has never read,
times node_shards, since shards on one node fill one directory. Over
bypass_bytes it is :bypassed; otherwise :used.
- The reply carries the shard's decision. The job reports the combined
decision (FileCache.combined/1): :used when any shard read through
its cache. job.scatter carries it too.
- smolquery_query_file_cache_shards_total{decision} counts each shard's
decision. smolquery_query_file_cache_jobs_total now counts once the
job has run, so a scattered job counts by what its shards did.
Rolling upgrade: the request still carries the coordinator's verdict as
file_cache, so a shard on an older node behaves as before, and a new
shard that gets no mode from an older coordinator follows that verdict.
Tests: shard_decision against a swept index (warm, cold, hot, forced,
no cache, node_shards); combined; scatter integration with a cold cache
(every shard bypasses), a warm cache (every shard uses it, and false
turns it off), and a shard called with a stale :bypassed verdict that
decides :used from its own index.
chasers
added this pull request to stack #456
October 4, 2026 01:45
Code review of PR 454 (/code-review high). Fixed: - A shard multiplied its own cold bytes by node_shards. Round-robin shards differ in size, so shards on one node could split between used and bypassed, and a table under bypass_bytes in total could bypass, which a single-engine run never does. The request now carries node_files, the sealed files of every shard on that node, and shard_decision/3 sums their cold bytes, so every shard on a node decides alike and a small table always uses the cache. node_files is built once per job, not counted per shard. - A shard on a pre-T-630 node replies without file_cache. It followed the coordinator's verdict, so it now counts as that verdict, not nil. - The jobs metric fired only at the end of a successful execute, so a job that failed after planning went uncounted. It now fires for done and failed jobs alike. A scattered job's event no longer carries the coordinator's uncached_bytes, which came from another node's index. - The shard metric fired before the shard ran, and not at all on the old-coordinator path. It now fires once the partial answers, on both. - job.scatter no longer repeats file_cache; the runner pops it off. - The cache block name lives in one test helper, Smolquery.Test.FileCacheFixture, instead of three copies. Not changed: the planner's segment_files call for a scattered job. The coordinator's verdict still sets its own engine and the single-engine fallback, so the call is not wasted. Tests: shard_decision/3 over node files; a cold cache under the threshold is used by every shard; a job that fails after planning counts its decision.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
TL;DR: Each scatter shard now decides the file cache against its own node's cache. Before, the coordinator decided for every shard, so a warmed table bypassed the cache on every run.
Tracker: T-630. Follow-up to T-629 (PR 453).
The bug
:bypassed) to every shard. Shards with a warm cache skip it too.fileCache: truewarming run fills the peers' caches, but the next auto run bypasses again.What changed
node_files: the sealed files (url and bytes) of every shard on that node.node_filesagainst its own node's index (FileCache.shard_decision/3).FileCache.combined/1):usedwhen any shard read through its cache.Metrics
smolquery_query_file_cache_shards_total{decision}smolquery_query_file_cache_jobs_total{decision}Rolling upgrade
file_cache.How to review
lib/smolquery/query_service/file_cache.ex:shard_decision/3andcombined/1.lib/smolquery/query_service/partial_worker.ex:cache_decision/2.lib/smolquery/query_service/scatter.ex:node_files/1, request keys, gathered decisions.lib/smolquery/query_service/runner.ex:count_cache_decision/3andscatter_cache/2.Tests
shard_decision/3against a swept index (warm, cold, empty, forced, no cache).combined/1.file_cache: false→ every shard is off.:bypassedverdict and decides:usedfrom its own index. This test fails without the fix.Checks
mix precommit,mix ci,mix dialyzerpass locally.mix test --include integrationfor the scatter and client suites.Review
/code-review highfound 10 candidates. Fixed in "Review of T-630":node_files.null. Now they count as the coordinator's verdict.job.scatterno longer repeatsfile_cache.Smolquery.Test.FileCacheFixture).segment_filescall. The coordinator still needs its verdict for its own engine and the fallback.Watch out
commit_conflict(the sqlite catalog is locked while engines attach). It happened in 2 of 6 local runs, including in a test this PR does not touch. CI passed.The test cluster is one node, so the tests use local workers. The cross-node case is proven by the stale-verdict test, not by a real peer.
Shards balance by file index, not by file. A change in membership moves files between nodes, so some cache misses still happen after a roll.
🤖 Generated with Claude Code