From 1e2cd973f9bc4292edc619262b78d2dee9937e58 Mon Sep 17 00:00:00 2001 From: AlinsRan Date: Mon, 28 Sep 2026 08:56:16 +0800 Subject: [PATCH] fix(keys): keep a returning metric on its slot, and stop broadcasting A metric that expired and came back consumed a new slot number and bumped delete_count, which sends every worker through a full sync of key_count slots -- on the request path, since sync() is the first thing add() does and add() runs on every observation of a metric with an exptime. Measured on the shape of a gateway pod (10 workers, 512m dict, the three exporter.lua metrics, 140k label combinations, 200 req/s): one such return took the worker set from 14.9% to 30.2% CPU, and 50 of them over ten seconds put three workers at 98.6%, 93.3% and 87.9% of a core. Three changes, and none of them adds shared state: ttl() tells apart the two states get() reports alike. A slot past its ttl whose node is still in the dict is the key's own -- expire() resurrects it with its value intact -- so the reclaim round and sync_range leave it alone, and the key comes back on the same number. Listing it meanwhile costs nothing: metric_data() skips a key whose value has expired. Only a slot whose node is gone gives its number up. On its own this is most of the bloat: 100 metrics expiring and returning used to take 100 new numbers, or 2 when nothing disturbed the workers; now they take 0. Renewal reads its own slot and renews it there, instead of starting with sync(). That is one dict read where sync() is two, and it touches none of the shared counters, so an external bump cannot make the request path scan: 15.7% against 287.7% under a storm of 50. The broadcast is gone. What it protected was not the duplicate metrics of #11934 -- the index check in list() does that -- but a peer dropping the index entry of a key that had since moved to another slot, which hid a live metric from that peer's scrape. Dropping an index entry only when it still points at the slot being given up fixes that locally, in one line, and a fresh slot is above every worker's self.last anyway, so their incremental sync finds it. Verified on 10 workers at 300k series with 1,500 new series/s: no duplicates, no live value missing from the output, counter values exactly as driven. --- CHANGELOG.md | 9 ++ prometheus_keys.lua | 114 ++++++++++++++++-------- prometheus_test.lua | 212 ++++++++++++++++++++++++++++++++------------ 3 files changed, 239 insertions(+), 96 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 21f35ea..4ea88a2 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,6 +10,15 @@ of changes. `flush_expired()` call holds the dict lock for a whole backlog, and let callers take the reclamation over with the `auto_flush_expired` option and `prometheus:flush_expired()` (#23). +- Stop consuming a new index slot, and stop broadcasting to the other workers, + when a metric that has expired comes back. A slot whose node is still in the + dict is resurrected in place, and only a slot whose node is gone gives its + number up -- `ttl()` tells the two apart where `get()` cannot. Renewing a + metric now reads its own slot instead of the shared counters, so nothing + another worker does can send it through a full sync on the request path. One + returning series used to take the worker set from 14.9% to 30.2% CPU and 50 of + them to 287.7%, with three workers at ~90-99% of a core (apache/apisix#13658, + apache/apisix#11934). Design and measurements: api7/rfcs#290. ## 1.0.0 diff --git a/prometheus_keys.lua b/prometheus_keys.lua index ee56683..2b95ddc 100644 --- a/prometheus_keys.lua +++ b/prometheus_keys.lua @@ -56,25 +56,34 @@ end -- check and remove expired keys function KeyIndex:remove_expired_keys() + -- Reclaiming first means the scan below sees the final state of every slot, + -- so a number whose node has just been reclaimed is given up in this round + -- rather than the next. Callers that schedule flush_expired() themselves -- + -- in a single process rather than in every worker -- turn this off with + -- auto_flush_expired. + if self.auto_flush_expired then + self:flush_expired() + end + for i, _ in pairs(self.expire_keys) do - -- Read i-th key. If it is nil or ttl is < 0, it means it was expired - local ttl, err = self.dict:ttl(self.key_prefix .. i) - if not (ttl and ttl >= 0 or err and err ~= "not found") then - if self.keys[i] then - self.index[self.keys[i]] = nil - self.keys[i] = nil - end - self.expire_keys[i] = nil + -- A slot is in one of three states, and only ttl() tells them apart -- + -- get() reports the last two alike, as nil: + -- + -- (1) live: ttl > 0 + -- (2) past its ttl, node still in the dict: ttl < 0 + -- (3) node physically reclaimed: "not found" + -- + -- In state (2) the node is still there and expire() resurrects it with its + -- value intact, so the slot is left exactly as it is: the key keeps its + -- number and comes back on it. Listing it in the meantime costs nothing -- + -- metric_data() skips a key whose value has expired. Only state (3) gives + -- the number up. + local _, err = self.dict:ttl(self.key_prefix .. i) + if err == "not found" then + self:forget_slot(i) end end - -- The loop above only drops worker-local references, so the expired entries - -- still have to be reclaimed from the dict itself. Callers that schedule - -- flush_expired() themselves -- in a single process rather than in every - -- worker -- turn this off with auto_flush_expired. - if self.auto_flush_expired then - self:flush_expired() - end end @@ -146,14 +155,35 @@ function KeyIndex:sync_range(first, last) end end elseif self.keys[i] then - self.index[self.keys[i]] = nil - self.keys[i] = nil - self.expire_keys[i] = nil + -- get() cannot tell a slot that is merely past its ttl from one whose + -- node is gone, and the two must not be treated alike: the first keeps + -- its number (see remove_expired_keys), only the second gives it up. + local _, err = self.dict:ttl(self.key_prefix .. i) + if err == "not found" then + self:forget_slot(i) + end end end self.last = last end + +-- Drops every local reference to a slot whose node is gone. +-- +-- self.keys (slot -> key) and self.index (key -> slot) are two views of one +-- mapping, and the index entry is only this slot's to drop while it still +-- points here. A key that has since moved to another slot keeps a live entry +-- there, and dropping it would hide that key from list() although its value is +-- in the dict -- which is what the delete_count broadcast used to paper over. +function KeyIndex:forget_slot(i) + local key = self.keys[i] + if key and self.index[key] == i then + self.index[key] = nil + end + self.keys[i] = nil + self.expire_keys[i] = nil +end + -- Returns array of all keys. function KeyIndex:list() self:sync() @@ -194,6 +224,24 @@ function KeyIndex:add(key_or_keys, err_msg_lru_eviction, exptime) local retried = false local repairs = 0 local repair_forcible = false + + -- The common case by far: this key has a slot and the slot still holds it, + -- so only its ttl has to be pushed out. Reading the slot costs one dict + -- read where sync() costs two, and none of the shared counters are touched + -- -- so nothing another worker does can turn this into a full sync. + local mine = self.index[key] + if mine and exptime and self.dict:get(self.key_prefix .. mine) == key then + local ok, err = self.dict:expire(self.key_prefix .. mine, exptime) + if ok then + goto renewed + end + if err ~= "not found" then + ngx.log(ngx.ERR, "failed to renew expire for key '", key, "': ", + tostring(err)) + goto renewed + end + end + while true do local N = self:sync() if self.index[key] ~= nil then @@ -203,22 +251,15 @@ function KeyIndex:add(key_or_keys, err_msg_lru_eviction, exptime) local ok, err = self.dict:expire(self.key_prefix .. self.index[key], exptime) if not ok then if err == "not found" then - -- The slot already expired in the shared dict. Drop the stale - -- local state and bump delete_count so other workers do a full - -- sync and reclaim the slot; without this the old slot lingers in - -- their local self.keys while the metric is re-added at a new slot, - -- desynchronizing the index and causing duplicate metric emission. - -- The dict slot is already gone (expire returned "not found"), so - -- there is no slot to clear here. - local idx = self.index[key] - self.index[key] = nil - self.keys[idx] = nil - self.expire_keys[idx] = nil - self.deleted = self.deleted + 1 - local _, incr_err, forcible = self.dict:incr(self.delete_count, 1, 0) - if incr_err or forcible then - return incr_err or err_msg_lru_eviction - end + -- The node is gone, so this number is no longer ours to renew. + -- The key takes a fresh slot below, which is above every worker's + -- self.last and is therefore picked up by their incremental sync + -- -- nothing has to be broadcast. Bumping delete_count here, as + -- this used to, sent every worker through a full sync of + -- key_count slots on its request path, and made each of them + -- forget every slot that was merely past its ttl, so those + -- metrics took new numbers too (apache/apisix#13658). + self:forget_slot(self.index[key]) expired = true else -- Unexpected expire error: the slot may still be live, so leave it @@ -295,6 +336,7 @@ function KeyIndex:add(key_or_keys, err_msg_lru_eviction, exptime) end retried = true end + ::renewed:: end end @@ -305,9 +347,7 @@ end function KeyIndex:remove(key, err_msg_lru_eviction) local i = self.index[key] if i then - self.index[key] = nil - self.keys[i] = nil - self.expire_keys[i] = nil + self:forget_slot(i) self.dict:set(self.key_prefix .. i, nil) self.deleted = self.deleted + 1 diff --git a/prometheus_test.lua b/prometheus_test.lua index 0bfe05d..4a59170 100644 --- a/prometheus_test.lua +++ b/prometheus_test.lua @@ -24,8 +24,9 @@ function SimpleDict:add(k, v, exptime) if k == "willnotfitk" or v == "willnotfitv" then forcible = true end - self:get(k) -- prunes the key if it has expired, like ngx.shared.DICT does - if self.dict and self.dict[k] then + -- ngx.shared.DICT:add only refuses a node that is still live; one that is + -- past its exptime is reused in place, keeping the same key. + if self:get(k) ~= nil then return false, "exists", false -- match ngx.shared.DICT:add on present keys end self:set(k, v, exptime) @@ -51,18 +52,27 @@ function SimpleDict:get(k) return nil, "dict error" end if not self.dict then self.dict = {} end - if self.dict[k] and self.dict[k]["expired"] and self.dict[k]["expired"] < os.time() then self.dict[k] = nil end - local v = self.dict[k] or {} - return v["value"], nil -- value, err + -- An entry past its exptime reads as missing but is NOT freed here: like + -- ngx.shared.DICT, only flush_expired() (or a write reusing the node) frees + -- it. Pruning here would collapse the "past its ttl" and "node reclaimed" + -- states, which ttl() has to tell apart. + local e = self.dict[k] + if not e or (e["expired"] and e["expired"] < os.time()) then + return nil, nil + end + return e["value"], nil -- value, err end function SimpleDict:delete(k) self.dict[k] = nil end function SimpleDict:expire(k, exptime) + -- Like ngx.shared.DICT:expire: it looks the node up without checking the + -- exptime, so a node that is merely past it is resurrected in place, with + -- its value intact. Only a node that has been freed reports "not found". if not self.dict[k] then return nil, "not found" end - self.dict[k]["expired"] = os.time() + exptime + self.dict[k]["expired"] = exptime ~= 0 and (os.time() + exptime) or nil return true -- match ngx.shared.DICT:expire, which returns true on success end function SimpleDict:flush_expired(n) @@ -82,18 +92,19 @@ function SimpleDict:flush_expired(n) return flushed end function SimpleDict:ttl(k) - -- Like ngx.shared.DICT:ttl, an expired entry reads as "not found" but is NOT - -- freed here: only flush_expired (or a write reusing the node) reclaims it. - -- Do not prune, or the physical-reclamation assertions below become vacuous. + -- Like ngx.shared.DICT:ttl, which peeks at the node without checking the + -- exptime. It reports the three states a slot can be in: + -- live -> a positive number (0 when permanent) + -- past its exptime, node kept -> a negative number + -- node freed -> nil, "not found" local e = self.dict and self.dict[k] - if not e or (e["expired"] and e["expired"] < os.time()) then + if not e then return nil, "not found" end - if e["expired"] then - return e["expired"] - os.time() - else + if not e["expired"] then return 0 end + return e["expired"] - os.time() end local function sleep(n) @@ -829,64 +840,131 @@ function TestKeyIndex:testList() luaunit.assertEquals(keys[3], "key3") end --- Regression test for apache/apisix#11934 (duplicate metrics). --- A key with an exptime is added and synced (self.last now tracks N). The key --- then expires in the underlying shared dict without going through remove(), so --- neither key_count nor delete_count changes and the next sync() is a no-op, --- leaving self.index still pointing at the now-vanished slot. Re-adding the key --- therefore takes the "expired" path in add() and allocates a new slot while --- the stale slot lingers in self.keys. Before the fix, delete_count was not --- bumped (other workers never re-synced) and list() iterated self.keys, so the --- same key was emitted twice -> duplicate metrics. -function TestKeyIndex:testExpiredReAddNoDuplicate() +-- A metric that expires and comes back while its node is still in the dict +-- keeps its slot: expire() resurrects the node with its value intact, so no +-- number is consumed and nothing has to be told to the other workers. +function TestKeyIndex:testExpiredReAddKeepsItsSlot() + -- reclaiming is left to the caller here, so the node stays in the dict and + -- the slot is in the state this test is about + self.key_index = require('prometheus_keys').new(self.dict, "_prefix_", 1, false) local err = self.key_index:add("expkey", "eviction_err", 1) luaunit.assertEquals(err, nil) self.key_index:sync() luaunit.assertEquals(self.dict:get("_prefix_key_count"), 1) luaunit.assertEquals(self.dict:get("_prefix_key_1"), "expkey") - luaunit.assertEquals(self.dict:get("_prefix_delete_count"), nil) - luaunit.assertEquals(self.key_index.index["expkey"], 1) - -- A second worker sharing the same shared dict syncs the initial state, so it - -- now holds slot 1 in its local self.keys/index. This is the worker that the - -- delete_count bump must later force to re-sync and reclaim the stale slot. + -- a second worker holds the same slot in its local state local worker2 = require('prometheus_keys').new(self.dict, "_prefix_", 1) worker2:sync() luaunit.assertEquals(worker2.index["expkey"], 1) - luaunit.assertEquals(#worker2:list(), 1) - -- Let the key expire in the underlying shared dict. key_count and - -- delete_count are untouched, so the next sync() inside add() is a no-op and - -- the local index keeps pointing at the (now gone) slot 1. sleep(2) + -- past its ttl, but the node is still there luaunit.assertEquals(self.dict:get("_prefix_key_1"), nil) + luaunit.assertTrue(self.dict:ttl("_prefix_key_1") < 0) + -- a reclaim round must leave the number alone in this state + self.key_index:remove_expired_keys() + luaunit.assertEquals(self.key_index.index["expkey"], 1) - -- Re-adding the now-expired key takes the expired branch and allocates slot 2. err = self.key_index:add("expkey", "eviction_err", 1) luaunit.assertEquals(err, nil) - luaunit.assertEquals(self.dict:get("_prefix_key_2"), "expkey") - - -- delete_count must have been bumped on the expired re-add path so that - -- other workers do a full sync and drop the stale slot. - luaunit.assertEquals(self.dict:get("_prefix_delete_count"), 1) - -- list() must report the key exactly once, not twice. - local keys = self.key_index:list() - luaunit.assertEquals(#keys, 1) - luaunit.assertEquals(keys[1], "expkey") + -- same slot, no new number, and nothing broadcast + luaunit.assertEquals(self.dict:get("_prefix_key_1"), "expkey") + luaunit.assertEquals(self.dict:get("_prefix_key_count"), 1) + luaunit.assertEquals(self.dict:get("_prefix_key_2"), nil) + luaunit.assertEquals(self.dict:get("_prefix_delete_count"), nil) - -- The second worker must converge: its next sync() sees the bumped - -- delete_count, does a full sync, drops the stale slot 1 and picks up slot 2. - -- Without the delete_count bump it would keep slot 1 forever and list() the - -- key twice. + -- both workers list it exactly once; worker2 never had to re-sync + luaunit.assertEquals(self.key_index:list(), {"expkey"}) + luaunit.assertEquals(worker2:list(), {"expkey"}) +end + + +-- Regression test for apache/apisix#11934 (duplicate metrics) and for removing +-- the delete_count broadcast that came with it. +-- +-- Once the node itself is gone the number is not the key's any more, so the key +-- takes a fresh slot. That slot is above every worker's self.last, so their +-- incremental sync picks it up: nothing has to be broadcast. The peer that +-- still holds the old, dead slot must neither list the key twice nor -- once +-- its own reclaim round runs -- lose it, which is what the broadcast used to +-- paper over. +function TestKeyIndex:testReclaimedSlotTakesANewOneWithoutBroadcasting() + luaunit.assertEquals(self.key_index:add("expkey", "eviction_err", 1), nil) + -- a handful of other metrics, so the dead slot sits below the peer's + -- incremental range, as an old slot on a long-running gateway does + for i = 1, 5 do + luaunit.assertEquals(self.key_index:add("other" .. i, "eviction_err"), nil) + end + local worker2 = require('prometheus_keys').new(self.dict, "_prefix_", 1) worker2:sync() - luaunit.assertEquals(worker2.keys[1], nil) - luaunit.assertEquals(worker2.index["expkey"], 2) - local keys2 = worker2:list() - luaunit.assertEquals(#keys2, 1) - luaunit.assertEquals(keys2[1], "expkey") + luaunit.assertEquals(worker2.index["expkey"], 1) + + sleep(2) + self.key_index:flush_expired() + luaunit.assertNil(self.dict.dict["_prefix_key_1"]) + + luaunit.assertEquals(self.key_index:add("expkey", "eviction_err", 60), nil) + local moved = self.key_index.index["expkey"] + luaunit.assertNotEquals(moved, 1) + luaunit.assertEquals(self.dict:get("_prefix_key_" .. moved), "expkey") + -- no broadcast: nothing forces the other workers into a full sync + luaunit.assertEquals(self.dict:get("_prefix_delete_count"), nil) + + -- the peer finds the new slot through its incremental sync, and lists the + -- key once although its own self.keys still holds the dead slot + local listed = {} + for _, k in ipairs(worker2:list()) do + listed[k] = (listed[k] or 0) + 1 + end + luaunit.assertEquals(listed["expkey"], 1) + + -- and its own reclaim round must not take the key down with the dead slot + worker2:remove_expired_keys() + luaunit.assertEquals(worker2.index["expkey"], moved) + listed = {} + for _, k in ipairs(worker2:list()) do + listed[k] = (listed[k] or 0) + 1 + end + luaunit.assertEquals(listed["expkey"], 1) +end + + +-- A live slot must survive a reclaim round untouched. +function TestKeyIndex:testRemoveExpiredKeysKeepsLiveSlots() + luaunit.assertEquals(self.key_index:add("livekey", "eviction_err", 60), nil) + luaunit.assertEquals(self.key_index:add("permkey", "eviction_err"), nil) + + self.key_index:remove_expired_keys() + + local listed = {} + for _, k in ipairs(self.key_index:list()) do + listed[k] = true + end + luaunit.assertTrue(listed["livekey"]) + luaunit.assertTrue(listed["permkey"]) + luaunit.assertEquals(self.key_index.index["livekey"], 1) +end + + +-- With no exptime -- the default in APISIX -- nothing ever expires, so none of +-- this comes into play and no shared counter is written at all. +function TestKeyIndex:testPermanentMetricsTouchNoSharedCounters() + for i = 1, 3 do + luaunit.assertEquals(self.key_index:add("perm" .. i, "eviction_err"), nil) + end + self.key_index:remove_expired_keys() + for i = 1, 3 do + luaunit.assertEquals(self.key_index:add("perm" .. i, "eviction_err"), nil) + end + + luaunit.assertEquals(self.dict:get("_prefix_key_count"), 3) + luaunit.assertEquals(self.dict:get("_prefix_delete_count"), nil) + luaunit.assertEquals(#self.key_index:list(), 3) end + -- remove_expired_keys() must physically reclaim the shared-dict space of -- expired entries, not just drop the worker-local references to them. -- Logically expired entries keep occupying slab pages, and the passive @@ -903,9 +981,11 @@ function TestKeyIndex:testRemoveExpiredKeysReclaimsSharedDict() sleep(2) - -- Both entries are logically gone but still physically present: every dict - -- API reports them as missing while they still hold their slab pages. - luaunit.assertEquals(self.dict:ttl("_prefix_key_1"), nil) + -- Both entries are logically gone but still physically present: get() + -- reports them as missing while they still hold their slab pages, and ttl() + -- is what tells that state apart from a node that is really gone. + luaunit.assertEquals(self.dict:get("_prefix_key_1"), nil) + luaunit.assertTrue(self.dict:ttl("_prefix_key_1") < 0) luaunit.assertNotNil(self.dict.dict["_prefix_key_1"]) luaunit.assertNotNil(self.dict.dict["expkey"]) @@ -918,6 +998,13 @@ function TestKeyIndex:testRemoveExpiredKeysReclaimsSharedDict() luaunit.assertNil(self.dict.dict["expkey"]) -- Entries that have not expired must be left alone. luaunit.assertNotNil(self.dict.dict["permanent"]) + + -- Only now, with the node gone, is the number given up: a second round sees + -- "not found" and drops the local reference to it. + local _, err2 = self.dict:ttl("_prefix_key_1") + luaunit.assertEquals(err2, "not found") + self.key_index:remove_expired_keys() + luaunit.assertNil(self.key_index.index["expkey"]) end -- flush_expired() holds the dict lock for its whole scan of the LRU queue, so @@ -967,7 +1054,9 @@ function TestKeyIndex:testAutoFlushExpiredDisabled() sleep(2) key_index:remove_expired_keys() - luaunit.assertNil(key_index.index["expkey"]) + -- the number is kept: the node is only past its ttl, so the key comes back + -- on the same slot + luaunit.assertEquals(key_index.index["expkey"], 1) -- both entries are still physically present, unlike with the default luaunit.assertNotNil(self.dict.dict["_noflush_key_1"]) luaunit.assertNotNil(self.dict.dict["expkey"]) @@ -1199,8 +1288,10 @@ function TestPrometheus:testKeyTimeout() self.p.key_index:sync() luaunit.assertEquals(self.dict:get("metric_exp"), nil) luaunit.assertEquals(self.dict:get("__ngx_prom__key_" .. i), nil) - luaunit.assertEquals(self.p.key_index.index["metric_exp"], nil) - luaunit.assertEquals(self.p.key_index.keys[i], nil) + -- the node is only past its ttl, so the slot is still this key's; the metric + -- is absent from the output because metric_data() skips a missing value + luaunit.assertEquals(self.p.key_index.index["metric_exp"], i) + luaunit.assertEquals(self.dict:get("metric_exp"), nil) self.gauge_exp:inc(1) luaunit.assertEquals(self.dict:get("gauge_exp"), 1) @@ -1212,8 +1303,11 @@ function TestPrometheus:testKeyTimeout() self.p.key_index:sync() luaunit.assertEquals(self.dict:get("gauge_exp"), nil) luaunit.assertEquals(self.dict:get("__ngx_prom__key_" .. i), nil) - luaunit.assertEquals(self.p.key_index.index["gauge_exp"], nil) - luaunit.assertEquals(self.p.key_index.keys[i], nil) + luaunit.assertEquals(self.p.key_index.index["gauge_exp"], i) + -- the reference survives as long as the node does; the metric is out of the + -- output because its value is gone, which is what metric_data() checks + luaunit.assertEquals(self.p.key_index.keys[i], "gauge_exp") + luaunit.assertEquals(self.dict:get("gauge_exp"), nil) self.gauge_exp_2:set(1) self.p.key_index:sync()