Skip to content

refactor: Introduce nebulex cache wrapper - #4054

Draft
bblaszkow06 wants to merge 6 commits into
Logflare:mainfrom
bblaszkow06:bb/o11y-1975-kv-nebulex-cache
Draft

bblaszkow06 wants to merge 6 commits into
Logflare:mainfrom
bblaszkow06:bb/o11y-1975-kv-nebulex-cache

Conversation

@bblaszkow06

Copy link
Copy Markdown
Contributor

This is meant to build a base for experimenting with a multi-layer Nebulex cache.

It was tested with the following benchmark:

Click to see the benchmark
# Usage: mix run bench/key_values_cache_bench.exs
#
# Measures the abstraction overhead of fronting Cachex with a Nebulex multi-level
# cache (the new `Logflare.KeyValues.Cache`) versus talking to a bare Cachex
# instance directly (the previous implementation). Compares the operations where
# the Nebulex dispatch is pure overhead with no DB call to mask it:
#
#   * read hit          - the ingestion hot path (LogEvent kv enrichment)
#   * negative-cache hit - a cached `nil` lookup
#   * write             - single entry put (multi-level write-through)
#   * bust_by/1         - WAL-driven invalidation over a set of lookup keys
#
# DB-miss (uncached) reads are intentionally excluded: DB latency dominates and
# would hide the wrapper overhead.

alias Logflare.KeyValues.Cache
alias Logflare.KeyValues.Cache.L1
alias Logflare.KeyValues.Cache.Multilevel
alias Logflare.Utils

# `mix run` boots the full application, which already starts the new Nebulex
# multi-level `Logflare.KeyValues.Cache` (and its Cachex-backed L1 level).

# Bare Cachex instance mirroring the previous KeyValues.Cache configuration, used
# as the baseline. Values are wrapped in `{:cached, _}` like the old read-through.
{:ok, _} =
  Cachex.start_link(:raw_kv_cache,
    compressed: true,
    hooks: [Utils.cache_limit(10_000_000)],
    expiration: Utils.cache_expiration_min(1440, 60)
  )

user_id = 1
value = %{"org_id" => "org_abc", "org" => %{"id" => "abc", "name" => "Acme"}}

hit_key = {:lookup, [user_id, "hit", nil]}
neg_key = {:lookup, [user_id, "neg", nil]}

# Previous-implementation read-through against the bare Cachex instance.
raw_lookup = fn key ->
  :raw_kv_cache
  |> Cachex.fetch(key, fn _ -> {:commit, {:cached, value}} end)
  |> case do
    {:commit, {:cached, v}} -> v
    {:ok, {:cached, v}} -> v
  end
end

# Previous-implementation bust for {user_id, key}: scan keys, filter, take.
raw_bust = fn uid, key ->
  {:ok, keys} = Cachex.keys(:raw_kv_cache)

  entries =
    Enum.filter(keys, fn
      {:lookup, [^uid, ^key | _]} -> true
      {:count, ^uid} -> true
      _ -> false
    end)

  Cachex.execute(:raw_kv_cache, fn worker ->
    Enum.reduce(entries, 0, fn k, acc ->
      case Cachex.take(worker, k) do
        {:ok, nil} -> acc
        {:ok, _} -> acc + 1
      end
    end)
  end)
end

# Underlying ETS table of the Cachex-backed L1 level (named by the adapter as
# `<cache>.Cachex`). Bust candidates can be matched here with a raw ETS match
# spec whose head bakes in user_id/key structurally (no interpreted guards).
l1_table = Logflare.KeyValues.Cache.L1.Cachex

# Structural ETS match: returns only the accessor variants for {uid, key},
# then deletes through Nebulex so the cache stays consistent.
ets_select_bust = fn uid, key ->
  spec = [{{:entry, {:lookup, [uid, key, :"$1"]}, :_, :_, :_}, [], [:"$1"]}]
  accessors = :ets.select(l1_table, spec)
  keys = [{:count, uid} | Enum.map(accessors, &{:lookup, [uid, key, &1]})]
  Multilevel.delete_all(in: keys)
end

# Speed ceiling: single-pass `:ets.select_delete` (bypasses Cachex/Nebulex).
ets_select_delete_bust = fn uid, key ->
  spec = [{{:entry, {:lookup, [uid, key, :_]}, :_, :_, :_}, [], [true]}]
  count = :ets.select_delete(l1_table, spec)
  Multilevel.delete({:count, uid})
  count
end

# Find via raw ETS structural match, delete at the L1 level (skips the multilevel
# delete_all, which appears to fall back to a full-table scan for key lists).
ets_select_l1_delete_bust = fn uid, key ->
  spec = [{{:entry, {:lookup, [uid, key, :"$1"]}, :_, :_, :_}, [], [:"$1"]}]
  accessors = :ets.select(l1_table, spec)
  keys = [{:count, uid} | Enum.map(accessors, &{:lookup, [uid, key, &1]})]
  L1.delete_all(in: keys)
end

# `delete_all(in: keys)` compiles every key into one match spec and runs
# `:ets.select_delete`, so its cost is (table size x key count). Deleting keys
# one by one is an O(1) hash lookup each. These variants isolate that.
ets_select_per_key_delete_bust = fn uid, key ->
  spec = [{{:entry, {:lookup, [uid, key, :"$1"]}, :_, :_, :_}, [], [:"$1"]}]
  accessors = :ets.select(l1_table, spec)
  keys = [{:count, uid} | Enum.map(accessors, &{:lookup, [uid, key, &1]})]
  Enum.each(keys, &Multilevel.delete/1)
end

ets_select_per_key_l1_delete_bust = fn uid, key ->
  spec = [{{:entry, {:lookup, [uid, key, :"$1"]}, :_, :_, :_}, [], [:"$1"]}]
  accessors = :ets.select(l1_table, spec)
  keys = [{:count, uid} | Enum.map(accessors, &{:lookup, [uid, key, &1]})]
  Enum.each(keys, &L1.delete/1)
end

# Same per-key delete, but the candidates come from the Nebulex stream API
# instead of a raw ETS match spec.
stream_per_key_delete_bust = fn uid, key ->
  keys =
    [select: :key]
    |> L1.stream!()
    |> Stream.filter(fn
      {:lookup, [^uid, ^key | _]} -> true
      _ -> false
    end)
    |> Enum.to_list()

  Enum.each([{:count, uid} | keys], &Multilevel.delete/1)
end

# `max_entries` is forwarded to `:ets.select/3` as the continuation buffer size
# and defaults to 100. Raising it is a pessimisation: building one big result
# list costs more than the continuation round-trips it saves.
stream_buffered_per_key_delete_bust = fn uid, key ->
  keys =
    [select: :key]
    |> L1.stream!(max_entries: 100_000)
    |> Stream.filter(fn
      {:lookup, [^uid, ^key | _]} -> true
      _ -> false
    end)
    |> Enum.to_list()

  Enum.each([{:count, uid} | keys], &Multilevel.delete/1)
end

# Find via Nebulex L1 stream + Elixir filter, delete at the L1 level.
stream_l1_delete_bust = fn uid, key ->
  keys =
    [select: :key]
    |> L1.stream!()
    |> Stream.filter(fn
      {:lookup, [^uid, ^key | _]} -> true
      _ -> false
    end)
    |> Enum.to_list()

  L1.delete_all(in: [{:count, uid} | keys])
end

# Background noise: unrelated lookup keys for other users. The bust path must
# scan past these, so this is what exposes the difference between filtering
# inside ETS (match spec) and copying the whole keyspace into the BEAM.
noise_count = 50_000

for i <- 1..noise_count do
  noise_key = {:lookup, [1_000_000 + i, "noise", nil]}
  Multilevel.put(noise_key, value)
  Cachex.put(:raw_kv_cache, noise_key, {:cached, value})
end

# Warm the read-hit and negative-cache entries on both caches (no DB during the
# measured runs).
Multilevel.put(hit_key, value)
Multilevel.put(neg_key, nil)
Cachex.put(:raw_kv_cache, hit_key, {:cached, value})
Cachex.put(:raw_kv_cache, neg_key, {:cached, nil})

# Pre-seed lookup keys for the bust benchmark (re-seeded inside the bench since
# bust deletes them).
seed_bust = fn ->
  for i <- 1..20 do
    Multilevel.put({:lookup, [user_id, "bust", "a#{i}"]}, value)
    Cachex.put(:raw_kv_cache, {:lookup, [user_id, "bust", "a#{i}"]}, {:cached, value})
  end
end

Benchee.run(
  %{
    "read hit (nebulex)" => fn -> Cache.lookup(user_id, "hit", nil) end,
    "read hit (raw cachex)" => fn -> raw_lookup.(hit_key) end
  },
  time: 2,
  warmup: 1,
  print: [fast_warning: false]
)

Benchee.run(
  %{
    "negative-cache hit (nebulex)" => fn -> Cache.lookup(user_id, "neg", nil) end,
    "negative-cache hit (raw cachex)" => fn -> raw_lookup.(neg_key) end
  },
  time: 2,
  warmup: 1,
  print: [fast_warning: false]
)

Benchee.run(
  %{
    "write (nebulex)" => fn -> Multilevel.put({:lookup, [user_id, "w", nil]}, value) end,
    "write (raw cachex)" => fn ->
      Cachex.put(:raw_kv_cache, {:lookup, [user_id, "w", nil]}, {:cached, value})
    end
  },
  time: 2,
  warmup: 1,
  print: [fast_warning: false]
)

Benchee.run(
  %{
    "bust_by (nebulex)" => {
      fn _ -> Cache.bust_by(user_id: user_id, key: "bust") end,
      before_each: fn _ -> seed_bust.() end
    },
    "bust_by (raw cachex)" => {
      fn _ -> raw_bust.(user_id, "bust") end,
      before_each: fn _ -> seed_bust.() end
    },
    "bust_by (ets select + multilevel delete)" => {
      fn _ -> ets_select_bust.(user_id, "bust") end,
      before_each: fn _ -> seed_bust.() end
    },
    "bust_by (ets select + L1 delete)" => {
      fn _ -> ets_select_l1_delete_bust.(user_id, "bust") end,
      before_each: fn _ -> seed_bust.() end
    },
    "bust_by (stream + L1 delete)" => {
      fn _ -> stream_l1_delete_bust.(user_id, "bust") end,
      before_each: fn _ -> seed_bust.() end
    },
    "bust_by (ets select_delete)" => {
      fn _ -> ets_select_delete_bust.(user_id, "bust") end,
      before_each: fn _ -> seed_bust.() end
    },
    "bust_by (ets select + per-key multilevel delete)" => {
      fn _ -> ets_select_per_key_delete_bust.(user_id, "bust") end,
      before_each: fn _ -> seed_bust.() end
    },
    "bust_by (ets select + per-key L1 delete)" => {
      fn _ -> ets_select_per_key_l1_delete_bust.(user_id, "bust") end,
      before_each: fn _ -> seed_bust.() end
    },
    "bust_by (stream + per-key multilevel delete)" => {
      fn _ -> stream_per_key_delete_bust.(user_id, "bust") end,
      before_each: fn _ -> seed_bust.() end
    },
    "bust_by (buffered stream + per-key multilevel delete)" => {
      fn _ -> stream_buffered_per_key_delete_bust.(user_id, "bust") end,
      before_each: fn _ -> seed_bust.() end
    }
  },
  time: 4,
  warmup: 4
)

Impact on hot path

op raw Cachex Nebulex delta
read hit 0.29 µs 2.54 µs +2.41 µs (8.3×)
negative-cache hit 0.21 µs 2.46 µs +2.39 µs (10.2×)
write 0.38 µs 1.58 µs +1.25 µs (3.9×)

The relative impact is high, but ~2µs absolutely

Invalidation

There was a 4x regression caused by faulty delete_all(in: keys) (now fixed in elixir-nebulex/nebulex_local#9 but unreleased)
But after the fix, it's faster than the original Cachex code

variant median
:ets.select_delete (raw ETS ceiling) 9.09 ms
ets select + per-key delete 9.69–10.22 ms
current impl 12.03 ms
old raw Cachex baseline 13.38 ms
buffered stream (max_entries: 100_000) 13.47 ms
any delete_all(in:) variant 50.5–52.0 ms

@bblaszkow06
bblaszkow06 force-pushed the bb/o11y-1975-kv-nebulex-cache branch 2 times, most recently from 7ceb7b6 to 57688ed Compare September 29, 2026 15:28
bblaszkow06 added a commit to bblaszkow06/logflare that referenced this pull request Sep 30, 2026
Generic code (health check, telemetry, test support, cache busting) called
Cachex directly on every cache module. That breaks as soon as a cache moves to
another backend, e.g. the Nebulex multilevel KeyValues cache in Logflare#4054.

- Logflare.Cache: a behaviour for cache operations (healthy?/0, stats/0,
  reset/0). `use Logflare.Cache` injects overridable Cachex defaults; the
  default implementation module can be swapped with `impl:`.
- `use Logflare.ContextCache` builds on it and makes bust_by/1 a required
  callback with a default primary-key bust. Every bust now goes through the
  cache's bust_by/1; the tuples CacheBuster emits are unchanged.
- The health check, telemetry and test resets call these callbacks instead of
  Cachex.
- Cachex configuration is shared (Logflare.Cache.CachexOps.child_spec/2 with
  named limit/ttl/purge_interval/warmer/compressed options), replacing the
  per-cache child_spec boilerplate and the Utils.cache_* helpers. Every cache
  keeps its current settings.

Behaviour changes:
- /health reports an unhealthy cache as "unhealthy" instead of "no_cache".
- The cachex.<cache>.purge and .stats metrics are removed; Cachex 4 never
  sets them, so they were always 0.
- Cache child specs are typed :supervisor (Cachex starts a supervisor).
- A {Rules, pkey} bust follows Rules.Cache's own `id:` busting instead of
  the generic scan; CacheBuster only sends Rules keyword busts.
- Busting a cache by an unsupported keyword raises ArgumentError.
bblaszkow06 added a commit to bblaszkow06/logflare that referenced this pull request Oct 1, 2026
Generic code (health check, telemetry, test support, cache busting) called
Cachex directly on every cache module. That breaks as soon as a cache moves to
another backend, e.g. the Nebulex multilevel KeyValues cache in Logflare#4054.

- Logflare.Cache: a behaviour for cache operations (healthy?/0, stats/0,
  reset/0). `use Logflare.Cache` injects overridable Cachex defaults; the
  default implementation module can be swapped with `impl:`.
- `use Logflare.ContextCache` builds on it and makes bust_by/1 a required
  callback with a default primary-key bust. Every bust now goes through the
  cache's bust_by/1; the tuples CacheBuster emits are unchanged.
- The health check, telemetry and test resets call these callbacks instead of
  Cachex.
- Cachex configuration is shared (Logflare.Cache.CachexOps.child_spec/2 with
  named limit/ttl/purge_interval/warmer/compressed options), replacing the
  per-cache child_spec boilerplate and the Utils.cache_* helpers. Every cache
  keeps its current settings.

Behaviour changes:
- /health reports an unhealthy cache as "unhealthy" instead of "no_cache".
- The cachex.<cache>.purge and .stats metrics are removed; Cachex 4 never
  sets them, so they were always 0.
- Cache child specs are typed :supervisor (Cachex starts a supervisor).
- A {Rules, pkey} bust follows Rules.Cache's own `id:` busting instead of
  the generic scan; CacheBuster only sends Rules keyword busts.
- Busting a cache by an unsupported keyword raises ArgumentError.
bblaszkow06 and others added 6 commits October 7, 2026 17:09
Generic code (health check, telemetry, test support, cache busting) called
Cachex directly on every cache module. That breaks as soon as a cache moves to
another backend, e.g. the Nebulex multilevel KeyValues cache in Logflare#4054.

- Logflare.Cache: a behaviour for cache operations (healthy?/0, stats/0,
  reset/0). `use Logflare.Cache` injects overridable Cachex defaults; the
  default implementation module can be swapped with `impl:`.
- `use Logflare.ContextCache` builds on it and makes bust_by/1 a required
  callback with a default primary-key bust. Every bust now goes through the
  cache's bust_by/1; the tuples CacheBuster emits are unchanged.
- The health check, telemetry and test resets call these callbacks instead of
  Cachex.
- Cachex configuration is shared (Logflare.Cache.CachexOps.child_spec/2 with
  named limit/ttl/purge_interval/warmer/compressed options), replacing the
  per-cache child_spec boilerplate and the Utils.cache_* helpers. Every cache
  keeps its current settings.

Behaviour changes:
- /health reports an unhealthy cache as "unhealthy" instead of "no_cache".
- The cachex.<cache>.purge and .stats metrics are removed; Cachex 4 never
  sets them, so they were always 0.
- Cache child specs are typed :supervisor (Cachex starts a supervisor).
- A {Rules, pkey} bust follows Rules.Cache's own `id:` busting instead of
  the generic scan; CacheBuster only sends Rules keyword busts.
- Busting a cache by an unsupported keyword raises ArgumentError.
Process.info(pid, :total_heap_size) returns the heap size in words, but the
cachex.<cache>.total_heap_size metric is declared with unit {:byte, :megabyte},
so the exported value was too small by a factor of the word size (8x on
64-bit). Convert words to bytes with :erlang.system_info(:wordsize) before
emitting it.
…ntation

Logflare.ContextCache still called Cachex directly for the read path
(apply_fun/3, update/4) and for the default primary-key bust.

- ContextCache gains required fetch/2 and update/2 callbacks next to
  bust_by/1. `use Logflare.ContextCache` injects overridable defaults that
  delegate to the `impl:` module (Logflare.Cache.CachexOps by default), and
  apply_fun/3, update/4 and bust_keys/1 only call these callbacks.
- The impl contract is explicit: Logflare.Cache.Ops (healthy?/1, stats/1,
  reset/1) and Logflare.ContextCache.Ops (bust_by/2, fetch/3, update/3).
- The Cachex read path, the primary-key bust and the gossip of cache misses
  move to CachexOps. The `{:cached, value}` wrapper is now private to it.
- ContextCache.fetch/3 and ContextCache.bust_by/2 are removed; callers use
  CachexOps directly.

Still Cachex-only: ContextCache.Gossip's receive side, cache warmers, and
cache modules that call Cachex themselves.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
TODO: Address the performance impact on busting
@bblaszkow06
bblaszkow06 force-pushed the bb/o11y-1975-kv-nebulex-cache branch from 57688ed to 39a2e6b Compare October 7, 2026 16:34

This branch has not been deployed

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

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant