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()