This is an automated email from the ASF dual-hosted git repository.
hanahmily pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/skywalking-banyandb.git
The following commit(s) were added to refs/heads/main by this push:
new b1fc90553 Add TTL to top-k query and add health check for pub (#1221)
b1fc90553 is described below
commit b1fc90553296510ec787c755411abe83ac12527e
Author: mrproliu <[email protected]>
AuthorDate: Wed Jul 22 06:43:51 2026 +0800
Add TTL to top-k query and add health check for pub (#1221)
* Add TTL to top-k query and add health check for pub
---
CHANGES.md | 2 +
banyand/liaison/grpc/bydbql.go | 33 +++++++---
banyand/liaison/grpc/bydbql_reparse_test.go | 2 +-
banyand/liaison/grpc/bydbql_test.go | 21 +++++-
banyand/liaison/grpc/server.go | 8 ++-
banyand/liaison/grpc/topk.go | 64 ++++++++++++++----
banyand/liaison/grpc/topk_test.go | 95 ++++++++++++++++++++++++---
banyand/queue/pub/health_check_wiring_test.go | 75 +++++++++++++++++++++
banyand/queue/pub/pub.go | 62 ++++++++++-------
docs/operation/configuration.md | 16 ++++-
10 files changed, 323 insertions(+), 55 deletions(-)
diff --git a/CHANGES.md b/CHANGES.md
index 95dbdbf6d..f78e52aaf 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -79,6 +79,7 @@ Release Notes.
- Add a Claude/Codex plugin packaging the BanyanDB MCP server with a
parse-only `validate_bydbql` tool (backed by the `mcp/tools/bydbql-parse` Go
validator) and a BydbQL skill for read-only natural-language-to-BydbQL
generation over STREAM/MEASURE/TRACE/PROPERTY resources.
- Introduce positional parameter binding (`?` placeholders) into BydbQL to
eliminate QL injection.
- Add reusable BydbQL binding: Prepare a query once, then Bind it many times
without re-parsing or mutating the template. The liaison caches prepared
statements on the gRPC query path (LRU bounded by entry count and bytes, on by
default; `--bydbql-prepared-cache-size`/`--bydbql-prepared-cache-max-bytes`,
`bydbql_prepared_cache_*` metrics) so repeated templates skip parsing. To
pinpoint un-cacheable and slow queries without high-cardinality labels, the
query access log tags each entry by [...]
+- Expire BydbQL top-K entries not seen within their TTL
(`--bydbql-topk-slow-ttl` / `--bydbql-topk-reparse-ttl`, default `24h`), and
log `last_seen` / `max_latency_at`.
### Bug Fixes
@@ -138,6 +139,7 @@ Release Notes.
- Deleting one TopN aggregation no longer tears down sibling aggregations on
the same source measure.
- Purge a deleted group's resource, index-rule and binding cache entries to
avoid dangling references.
- Clear a trace subject's index when its last index rule or binding is removed.
+- Enable periodic health checks on the queue client
(`--<prefix>-client-health-check-interval`, default `10s`), evicting dead data
nodes proactively.
### Document
diff --git a/banyand/liaison/grpc/bydbql.go b/banyand/liaison/grpc/bydbql.go
index 7a6422f2f..2f6496681 100644
--- a/banyand/liaison/grpc/bydbql.go
+++ b/banyand/liaison/grpc/bydbql.go
@@ -176,9 +176,11 @@ func (b *bydbQLService) Query(ctx context.Context, req
*bydbqlv1.QueryRequest) (
return resp, nil
}
-// topKDumper tracks the top re-parsed and slow queries and, on a supervised
-// goroutine, periodically logs the cumulative top-K. All methods are
nil-safe, so the
-// call sites need no guards when the top-K log is disabled (the dumper is
nil).
+// topKDumper tracks the top re-parsed and slow queries and, on a supervised
goroutine,
+// periodically logs the top-K, having first dropped whatever each tracker has
not seen
+// within its TTL. Entries expire by inactivity only: a query that keeps
recurring keeps its
+// accumulated count and peak. All methods are nil-safe, so the call sites
need no guards
+// when the top-K log is disabled (the dumper is nil).
type topKDumper struct {
// reparse holds only templates the cache had already compiled once and
had to
// compile again. First-ever compiles are excluded at the call site,
which is what
@@ -191,13 +193,19 @@ type topKDumper struct {
}
// newTopKDumper starts the trackers and the dump goroutine; a non-positive
interval
-// disables the feature and returns nil.
-func newTopKDumper(interval time.Duration, l *logger.Logger) *topKDumper {
+// disables the feature and returns nil. Each tracker expires entries on its
own TTL:
+// a template nobody has re-parsed, or a query nobody has run slowly, for that
long stops
+// being reported. A non-positive TTL keeps that tracker's entries for the
process lifetime.
+func newTopKDumper(interval, reparseTTL, slowTTL time.Duration, l
*logger.Logger) *topKDumper {
if interval <= 0 {
return nil
}
ctx, cancel := context.WithCancel(context.Background())
- d := &topKDumper{reparse: newTopK(bydbqlTopKSize), slow:
newTopK(bydbqlTopKSize), l: l, cancel: cancel}
+ d := &topKDumper{
+ reparse: newTopK(bydbqlTopKSize, reparseTTL, nil),
+ slow: newTopK(bydbqlTopKSize, slowTTL, nil),
+ l: l, cancel: cancel,
+ }
run.Go(ctx, "liaison.grpc.bydbql.topk-dump", l, func(ctx
context.Context) {
ticker := time.NewTicker(interval)
defer ticker.Stop()
@@ -245,13 +253,22 @@ func (d *topKDumper) dump() {
// root; keeping the threshold on top of it would only re-hide the real
re-parses the
// exclusion just made visible.
d.logTopK(d.reparse.snapshot(), 1, "top bydbql cache-miss queries",
func(s topKSlot) string {
- return fmt.Sprintf("%q count=%d", s.key, s.count)
+ return fmt.Sprintf("%q count=%d last_seen=%s", s.key, s.count,
formatTopKTime(s.lastSeen))
})
+ // max_latency is a running peak, so it can outlive the condition that
caused it by
+ // as long as the TTL allows. max_latency_at dates it, and last_seen
says when the
+ // query last ran at all, which together separate a live problem from a
stale peak.
d.logTopK(d.slow.snapshotByLatency(), 1, "top bydbql slow queries",
func(s topKSlot) string {
- return fmt.Sprintf("%q count=%d max_latency=%s", s.key,
s.count, s.maxDur)
+ return fmt.Sprintf("%q count=%d max_latency=%s
max_latency_at=%s last_seen=%s",
+ s.key, s.count, s.maxDur, formatTopKTime(s.maxDurAt),
formatTopKTime(s.lastSeen))
})
}
+// formatTopKTime renders a tracker timestamp in the same layout the
surrounding logs use.
+func formatTopKTime(t time.Time) string {
+ return t.UTC().Format(time.RFC3339)
+}
+
// logTopK logs the entries with at least minCount occurrences, formatted by
line.
func (d *topKDumper) logTopK(entries []topKSlot, minCount uint64, msg string,
line func(topKSlot) string) {
lines := formatTopK(entries, minCount, line)
diff --git a/banyand/liaison/grpc/bydbql_reparse_test.go
b/banyand/liaison/grpc/bydbql_reparse_test.go
index a17675968..e678389c2 100644
--- a/banyand/liaison/grpc/bydbql_reparse_test.go
+++ b/banyand/liaison/grpc/bydbql_reparse_test.go
@@ -31,7 +31,7 @@ import (
// change that broke it. run() returns the cache result for assertions.
func newReparseProbe(t *testing.T, size, maxBytes int) (*preparedCache, *topK,
func(query string) string) {
t.Helper()
- tk := newTopK(bydbqlTopKSize)
+ tk := newTopK(bydbqlTopKSize, 0, nil)
c := newPreparedCache(size, maxBytes, nil)
run := func(q string) string {
_, result, err := c.getOrPrepare(q)
diff --git a/banyand/liaison/grpc/bydbql_test.go
b/banyand/liaison/grpc/bydbql_test.go
index ae92bb8a6..a9ecb4bf3 100644
--- a/banyand/liaison/grpc/bydbql_test.go
+++ b/banyand/liaison/grpc/bydbql_test.go
@@ -95,7 +95,7 @@ func
TestBydbQLQuery_ParamTypeMismatch_ReturnsInvalidArgument(t *testing.T) {
// newTestDumper builds a topKDumper without starting the dump goroutine.
func newTestDumper(l *logger.Logger) *topKDumper {
- return &topKDumper{reparse: newTopK(bydbqlTopKSize), slow:
newTopK(bydbqlTopKSize), l: l}
+ return &topKDumper{reparse: newTopK(bydbqlTopKSize, 0, nil), slow:
newTopK(bydbqlTopKSize, 0, nil), l: l}
}
// attachTestDumper gives svc a dumper without starting its dump goroutine, so
the
@@ -157,6 +157,25 @@ func TestBydbQLDumpTopK(t *testing.T) {
assert.NotEmpty(t, d.slow.snapshot())
}
+// newTopKDumper takes the two TTLs as adjacent time.Duration parameters, so
swapping them
+// would compile and behave plausibly while silently applying each tracker's
TTL to the
+// other. Pin each one to its own tracker. The other tests build the dumper
struct directly,
+// so this is the only coverage of the real constructor.
+func TestNewTopKDumperWiresEachTTLToItsTracker(t *testing.T) {
+ const reparseTTL, slowTTL = 3 * time.Hour, 7 * time.Hour
+ d := newTopKDumper(time.Hour, reparseTTL, slowTTL,
logger.GetLogger("test-bydbql"))
+ require.NotNil(t, d)
+ defer d.close()
+
+ assert.Equal(t, reparseTTL, d.reparse.ttl, "the reparse tracker must
get the reparse TTL")
+ assert.Equal(t, slowTTL, d.slow.ttl, "the slow tracker must get the
slow TTL")
+}
+
+func TestNewTopKDumperDisabledByNonPositiveInterval(t *testing.T) {
+ assert.Nil(t, newTopKDumper(0, time.Hour, time.Hour,
logger.GetLogger("test-bydbql")),
+ "a non-positive interval disables the feature, and the nil
dumper's observers no-op")
+}
+
func TestFormatTopKAppliesItsMinCount(t *testing.T) {
entries := []topKSlot{
{key: "frequent", count: 5},
diff --git a/banyand/liaison/grpc/server.go b/banyand/liaison/grpc/server.go
index a9425a632..71af0bb9a 100644
--- a/banyand/liaison/grpc/server.go
+++ b/banyand/liaison/grpc/server.go
@@ -139,6 +139,8 @@ type server struct {
grpcBufferMemoryRatio float64
bydbqlSlowThreshold time.Duration
bydbqlTopKLogInterval time.Duration
+ bydbqlTopKSlowTTL time.Duration
+ bydbqlTopKReparseTTL time.Duration
bydbqlCacheSize int
bydbqlCacheMaxBytes int
port uint32
@@ -372,7 +374,7 @@ func (s *server) PreRun(ctx context.Context) error {
s.bydbQLSVC.slowThreshold = s.bydbqlSlowThreshold
// The dump goroutine lives for the server's lifetime and is stopped by
Close(),
// so it is rooted at a background context rather than PreRun's setup
context.
- s.bydbQLSVC.dumper = newTopKDumper(s.bydbqlTopKLogInterval,
s.bydbQLSVC.l) //nolint:contextcheck
+ s.bydbQLSVC.dumper = newTopKDumper(s.bydbqlTopKLogInterval,
s.bydbqlTopKReparseTTL, s.bydbqlTopKSlowTTL, s.bydbQLSVC.l)
//nolint:contextcheck
s.propertyServer.metrics = metrics
if s.barrierSVC != nil {
s.barrierSVC.metrics = metrics
@@ -466,6 +468,10 @@ func (s *server) FlagSet() *run.FlagSet {
"end-to-end latency above which a BydbQL query is counted as
slow; 0 disables slow-query tracking")
fs.DurationVar(&s.bydbqlTopKLogInterval, "bydbql-topk-log-interval",
5*time.Minute,
"how often to log the top BydbQL cache-miss and slow queries; 0
disables the top-K log")
+ fs.DurationVar(&s.bydbqlTopKSlowTTL, "bydbql-topk-slow-ttl",
24*time.Hour,
+ "drop a slow-query top-K entry whose query has not been slow
again for this long; 0 keeps it for the process lifetime")
+ fs.DurationVar(&s.bydbqlTopKReparseTTL, "bydbql-topk-reparse-ttl",
24*time.Hour,
+ "drop a cache-miss top-K entry whose template has not been
re-parsed again for this long; 0 keeps it for the process lifetime")
s.grpcBufferMemoryRatio = 0.1
fs.Float64Var(&s.grpcBufferMemoryRatio, "grpc-buffer-memory-ratio", 0.1,
"ratio of memory limit to use for gRPC buffer size calculation
(0.0 < ratio <= 1.0)")
diff --git a/banyand/liaison/grpc/topk.go b/banyand/liaison/grpc/topk.go
index f6ef17773..c1c3c724d 100644
--- a/banyand/liaison/grpc/topk.go
+++ b/banyand/liaison/grpc/topk.go
@@ -30,11 +30,15 @@ import (
// a normal workload, keeping the reported counts exact while still bounding
memory.
const bydbqlTopKSize = 128
-// topKSlot is one tracked query and its accumulated statistics.
+// topKSlot is one tracked query and its accumulated statistics. lastSeen
drives TTL
+// expiry; maxDurAt pins when the peak latency happened, so a consumer can
tell a live
+// problem from a peak the tracker has merely been carrying since a startup
incident.
type topKSlot struct {
- key string
- count uint64
- maxDur time.Duration
+ lastSeen time.Time
+ maxDurAt time.Time
+ key string
+ count uint64
+ maxDur time.Duration
}
// topK is a bounded approximate heavy-hitters tracker (Space-Saving): it
keeps at
@@ -42,17 +46,30 @@ type topKSlot struct {
// inherit that entry's count + 1 so a fresh key gets a fair chance instead of
being
// evicted again immediately. With k modest (128) the min scan and the
snapshot sort are
// cheap, and observing an existing key needs no reordering at all.
+//
+// Entries also expire, by inactivity: a key not observed for ttl is dropped,
so a query
+// that stops happening stops being reported. A key that does keep being
observed never
+// expires and goes on accumulating, so expiry bounds how long a finished
problem lingers,
+// not how far back an ongoing one is counted. A ttl <= 0 keeps entries for
the process
+// lifetime.
type topK struct {
slots map[string]*topKSlot
+ now func() time.Time
k int
+ ttl time.Duration
mu sync.Mutex
}
-func newTopK(k int) *topK {
+// newTopK builds a tracker holding at most k entries, expiring any entry not
observed
+// for ttl. now supplies the clock; passing nil uses time.Now.
+func newTopK(k int, ttl time.Duration, now func() time.Time) *topK {
if k < 1 {
k = 1
}
- return &topK{slots: make(map[string]*topKSlot, k), k: k}
+ if now == nil {
+ now = time.Now
+ }
+ return &topK{slots: make(map[string]*topKSlot, k), k: k, ttl: ttl, now:
now}
}
// observe records one occurrence of key; dur is the query latency (0 when
latency is
@@ -60,15 +77,22 @@ func newTopK(k int) *topK {
func (t *topK) observe(key string, dur time.Duration) {
t.mu.Lock()
defer t.mu.Unlock()
+ now := t.now()
if s, ok := t.slots[key]; ok {
s.count++
+ s.lastSeen = now
if dur > s.maxDur {
- s.maxDur = dur
+ s.maxDur, s.maxDurAt = dur, now
}
return
}
+ if len(t.slots) >= t.k {
+ // Reclaim expired entries before falling back to evicting a
live one: a key
+ // nobody has asked about in ttl is a better victim than a
merely infrequent one.
+ t.purgeExpiredLocked(now)
+ }
if len(t.slots) < t.k {
- t.slots[key] = &topKSlot{key: key, count: 1, maxDur: dur}
+ t.slots[key] = &topKSlot{key: key, count: 1, maxDur: dur,
lastSeen: now, maxDurAt: now}
return
}
// Full: evict the least-frequent entry and let the new key inherit its
count.
@@ -80,12 +104,26 @@ func (t *topK) observe(key string, dur time.Duration) {
}
}
delete(t.slots, minKey)
- t.slots[key] = &topKSlot{key: key, count: minCount + 1, maxDur: dur}
+ t.slots[key] = &topKSlot{key: key, count: minCount + 1, maxDur: dur,
lastSeen: now, maxDurAt: now}
+}
+
+// purgeExpiredLocked drops every entry not observed within ttl. Call with the
lock held.
+func (t *topK) purgeExpiredLocked(now time.Time) {
+ if t.ttl <= 0 {
+ return
+ }
+ for k, s := range t.slots {
+ if now.Sub(s.lastSeen) > t.ttl {
+ delete(t.slots, k)
+ }
+ }
}
// snapshot returns the tracked entries ranked by frequency: (count desc,
maxDur desc,
-// key asc). The tracker is cumulative, so each dump reflects the hottest
queries since
-// process start. The full tie-break makes the order deterministic across
dumps.
+// key asc). The TTL expires an entry by inactivity; it does not window the
statistics. For
+// as long as a key keeps being observed its count and maxDur stay cumulative,
so a
+// still-active entry can report a peak from hours ago — which is what
maxDurAt and lastSeen
+// exist to date. The full tie-break makes the order deterministic across
dumps.
func (t *topK) snapshot() []topKSlot {
out := t.copyOut()
sort.Slice(out, func(i, j int) bool { return lessByCount(out[i],
out[j]) })
@@ -101,8 +139,12 @@ func (t *topK) snapshotByLatency() []topKSlot {
return out
}
+// copyOut snapshots the live entries, expiring stale ones first. Purging here
rather
+// than only in observe is what lets a tracker that has gone quiet drain: an
entry no
+// one observes again would otherwise never be revisited, and would be
reported forever.
func (t *topK) copyOut() []topKSlot {
t.mu.Lock()
+ t.purgeExpiredLocked(t.now())
out := make([]topKSlot, 0, len(t.slots))
for _, s := range t.slots {
out = append(out, *s)
diff --git a/banyand/liaison/grpc/topk_test.go
b/banyand/liaison/grpc/topk_test.go
index db864ee03..8c3cd3718 100644
--- a/banyand/liaison/grpc/topk_test.go
+++ b/banyand/liaison/grpc/topk_test.go
@@ -24,6 +24,8 @@ import (
"time"
"github.com/stretchr/testify/assert"
+
+ "github.com/apache/skywalking-banyandb/pkg/timestamp"
)
func topKByKey(slots []topKSlot) map[string]topKSlot {
@@ -35,7 +37,7 @@ func topKByKey(slots []topKSlot) map[string]topKSlot {
}
func TestTopKCountsAndMaxDur(t *testing.T) {
- tk := newTopK(4)
+ tk := newTopK(4, 0, nil)
tk.observe("a", 5*time.Millisecond)
tk.observe("a", 8*time.Millisecond) // largest, observed in the middle
tk.observe("a", 3*time.Millisecond) // smaller and last: maxDur must
stay 8ms, not 3ms
@@ -48,7 +50,7 @@ func TestTopKCountsAndMaxDur(t *testing.T) {
}
func TestTopKEvictsMinAndInheritsCount(t *testing.T) {
- tk := newTopK(2)
+ tk := newTopK(2, 0, nil)
tk.observe("a", 0)
tk.observe("a", 0)
tk.observe("a", 0) // a.count = 3
@@ -67,7 +69,7 @@ func TestTopKEvictsMinAndInheritsCount(t *testing.T) {
}
func TestTopKSnapshotSortedAndCumulative(t *testing.T) {
- tk := newTopK(4)
+ tk := newTopK(4, 0, nil)
tk.observe("hi", 0)
tk.observe("hi", 0)
tk.observe("lo", 0)
@@ -80,7 +82,7 @@ func TestTopKSnapshotSortedAndCumulative(t *testing.T) {
}
func TestTopKSnapshotByLatency(t *testing.T) {
- tk := newTopK(bydbqlTopKSize)
+ tk := newTopK(bydbqlTopKSize, 0, nil)
tk.observe("frequent", time.Millisecond) // high count, low latency
tk.observe("frequent", time.Millisecond)
tk.observe("frequent", time.Millisecond)
@@ -92,9 +94,11 @@ func TestTopKSnapshotByLatency(t *testing.T) {
}
func TestTopKDeterministicTieBreak(t *testing.T) {
- // Equal counts and durations must order by key, not by map iteration.
- a := newTopK(bydbqlTopKSize)
- b := newTopK(bydbqlTopKSize)
+ // Equal counts and durations must order by key, not by map iteration.
Both trackers run
+ // on mock clocks frozen at the same instant, so the compared slots
differ only in order.
+ clockA, clockB := timestamp.NewMockClock(), timestamp.NewMockClock()
+ a := newTopK(bydbqlTopKSize, 0, clockA.Now)
+ b := newTopK(bydbqlTopKSize, 0, clockB.Now)
for _, k := range []string{"c", "a", "b"} {
a.observe(k, 0)
}
@@ -106,8 +110,83 @@ func TestTopKDeterministicTieBreak(t *testing.T) {
assert.Equal(t, []string{"a", "b", "c"}, []string{snap[0].key,
snap[1].key, snap[2].key})
}
+func TestTopKExpiresUnobservedEntry(t *testing.T) {
+ mc := timestamp.NewMockClock()
+ tk := newTopK(4, time.Hour, mc.Now)
+ tk.observe("stale", time.Second)
+
+ mc.Add(time.Hour - time.Minute)
+ assert.Len(t, tk.snapshot(), 1, "still within the TTL")
+
+ mc.Add(2 * time.Minute) // now past the TTL
+ assert.Empty(t, tk.snapshot(), "an entry not seen within the TTL is
dropped")
+}
+
+func TestTopKObserveRefreshesTTL(t *testing.T) {
+ mc := timestamp.NewMockClock()
+ tk := newTopK(4, time.Hour, mc.Now)
+ tk.observe("recurring", 0)
+
+ // Re-observing before expiry must reset the clock on the entry, so a
query that
+ // keeps happening is never dropped no matter how long the process has
run.
+ for i := 0; i < 5; i++ {
+ mc.Add(50 * time.Minute)
+ tk.observe("recurring", 0)
+ }
+ mc.Add(50 * time.Minute)
+
+ snap := tk.snapshot()
+ assert.Len(t, snap, 1, "a repeatedly observed entry survives
indefinitely")
+ assert.Equal(t, uint64(6), snap[0].count)
+}
+
+func TestTopKZeroTTLKeepsEntriesForever(t *testing.T) {
+ mc := timestamp.NewMockClock()
+ tk := newTopK(4, 0, mc.Now)
+ tk.observe("kept", 0)
+
+ mc.Add(365 * 24 * time.Hour)
+ assert.Len(t, tk.snapshot(), 1, "ttl <= 0 restores the cumulative
behavior")
+}
+
+func TestTopKPurgesExpiredBeforeEvictingLiveEntry(t *testing.T) {
+ mc := timestamp.NewMockClock()
+ tk := newTopK(2, time.Hour, mc.Now)
+ tk.observe("expiring", 0)
+ tk.observe("expiring", 0)
+ tk.observe("expiring", 0) // count 3: the most frequent, so never the
evict-min victim
+
+ mc.Add(2 * time.Hour) // "expiring" is now stale
+ tk.observe("fresh", 0) // fills the second slot
+ tk.observe("newcomer", 0)
+
+ byKey := topKByKey(tk.snapshot())
+ assert.Len(t, byKey, 2)
+ _, hasExpiring := byKey["expiring"]
+ assert.False(t, hasExpiring, "the stale entry is reclaimed even though
it had the highest count")
+ assert.Equal(t, uint64(1), byKey["newcomer"].count,
+ "reclaiming a slot lets the new key start at 1 instead of
inheriting an evicted count")
+}
+
+func TestTopKMaxDurAtTracksPeakNotLastObserve(t *testing.T) {
+ mc := timestamp.NewMockClock()
+ tk := newTopK(4, time.Hour, mc.Now)
+
+ tk.observe("q", time.Millisecond)
+ mc.Add(time.Minute)
+ tk.observe("q", time.Second) // the peak
+ peakAt := mc.Now()
+ mc.Add(time.Minute)
+ tk.observe("q", time.Millisecond) // slower observation must not move
maxDurAt
+
+ snap := tk.snapshot()
+ assert.Equal(t, time.Second, snap[0].maxDur)
+ assert.Equal(t, peakAt, snap[0].maxDurAt, "maxDurAt dates the peak, not
the latest observation")
+ assert.Equal(t, mc.Now(), snap[0].lastSeen, "lastSeen tracks the latest
observation")
+}
+
func TestTopKConcurrentObserve(t *testing.T) {
- tk := newTopK(bydbqlTopKSize)
+ tk := newTopK(bydbqlTopKSize, 0, nil)
const workers, iters = 16, 500
var wg sync.WaitGroup
for w := 0; w < workers; w++ {
diff --git a/banyand/queue/pub/health_check_wiring_test.go
b/banyand/queue/pub/health_check_wiring_test.go
new file mode 100644
index 000000000..c2fb017ba
--- /dev/null
+++ b/banyand/queue/pub/health_check_wiring_test.go
@@ -0,0 +1,75 @@
+// Licensed to Apache Software Foundation (ASF) under one or more contributor
+// license agreements. See the NOTICE file distributed with
+// this work for additional information regarding copyright
+// ownership. Apache Software Foundation (ASF) licenses this file to you under
+// the Apache License, Version 2.0 (the "License"); you may
+// not use this file except in compliance with the License.
+// You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+package pub
+
+import (
+ "context"
+ "testing"
+ "time"
+
+ "github.com/stretchr/testify/assert"
+ "github.com/stretchr/testify/require"
+ grpc_health_v1 "google.golang.org/grpc/health/grpc_health_v1"
+
+ databasev1
"github.com/apache/skywalking-banyandb/api/proto/banyandb/database/v1"
+ modelv1
"github.com/apache/skywalking-banyandb/api/proto/banyandb/model/v1"
+ "github.com/apache/skywalking-banyandb/pkg/logger"
+)
+
+// The periodic prober itself has long been implemented and tested in
pkg/grpchelper; what
+// was missing is that PreRun never passed an interval, so for the queue
client the prober
+// simply never started. The active set was then only ever corrected by a
request failing on
+// a dead node, which is why a node that disappeared stayed routable until
some query paid
+// the full timeout to discover it.
+//
+// These tests therefore pin the WIRING, not the mechanism. Both fail before
the fix: the
+// default assertion because no flag existed, and the eviction assertion
because a zero
+// interval leaves the prober unstarted, so the node stays active forever.
+func TestFlagSetSuppliesNonZeroHealthCheckInterval(t *testing.T) {
+ p := New(nil, databasev1.Role_ROLE_DATA).(*pub)
+ p.FlagSet() // DurationVar writes the default through immediately
+
+ assert.NotZero(t, p.healthCheckInterval,
+ "a zero interval silently disables the prober, which is the bug
this flag exists to prevent")
+ assert.Equal(t, defaultHealthCheckInterval, p.healthCheckInterval)
+}
+
+func TestPreRunStartsPeriodicHealthCheck(t *testing.T) {
+ addr := getAddress()
+ healthSrv, stopServer := setupWithStatus(addr,
modelv1.Status_STATUS_SUCCEED)
+ defer stopServer()
+
+ p := New(nil, databasev1.Role_ROLE_DATA).(*pub)
+ p.log = logger.GetLogger("queue-client")
+ p.FlagSet()
+ p.healthCheckInterval = 100 * time.Millisecond // keep the test short
+ require.NoError(t, p.PreRun(context.Background()))
+ defer p.connMgr.GracefulStop()
+
+ p.OnAddOrUpdate(getDataNode("node1", addr))
+ require.Eventually(t, func() bool { return p.connMgr.ActiveCount() == 1
},
+ 10*time.Second, 20*time.Millisecond, "the healthy node should
be admitted")
+
+ // The node stops serving. Nothing publishes to it, so only a
background prober can
+ // notice — which is exactly what a zero interval would fail to do.
+ healthSrv.SetServingStatus("",
grpc_health_v1.HealthCheckResponse_NOT_SERVING)
+
+ assert.Eventually(t, func() bool { return p.connMgr.ActiveCount() == 0
},
+ 10*time.Second, 20*time.Millisecond,
+ "the prober must evict the unreachable node without any request
being sent to it")
+}
diff --git a/banyand/queue/pub/pub.go b/banyand/queue/pub/pub.go
index a51996888..8c13a1729 100644
--- a/banyand/queue/pub/pub.go
+++ b/banyand/queue/pub/pub.go
@@ -54,6 +54,12 @@ import (
var queuePubScope = observability.RootScope.SubScope("queue_pub")
+// defaultHealthCheckInterval is how often the queue client re-checks the
nodes it considers
+// active, matching the property schema client's default. Without it the
active set is only
+// validated on admission and then by a request failing on it, so a dead node
keeps absorbing
+// queries — each paying the full broadcast timeout — until one of them evicts
it.
+const defaultHealthCheckInterval = 10 * time.Second
+
// ChunkedSyncClientConfig configures chunked sync client behavior.
type ChunkedSyncClientConfig struct {
ChunkSize uint32 // Size of each chunk in bytes
@@ -72,26 +78,27 @@ var (
type pub struct {
schema.UnimplementedOnInitHandler
- metadata metadata.Repo
- handlers map[bus.Topic]schema.EventHandler
- log *logger.Logger
- metrics *pubMetrics
- migrationMetrics *pubMigrationMetrics
- connMgr *grpchelper.ConnManager[*client]
- closer *run.Closer
- writableProbe map[string]map[string]struct{}
- nodeCache map[string]nodeInfo
- caCertPath string
- caCertReloader *pkgtls.Reloader
- prefix string
- retryPolicy string
- selfNode string
- selfRole string
- selfTier string
- allowedRoles []databasev1.Role
- writableProbeMu sync.Mutex
- nodeCacheMu sync.RWMutex
- tlsEnabled bool
+ metadata metadata.Repo
+ handlers map[bus.Topic]schema.EventHandler
+ log *logger.Logger
+ metrics *pubMetrics
+ migrationMetrics *pubMigrationMetrics
+ connMgr *grpchelper.ConnManager[*client]
+ closer *run.Closer
+ writableProbe map[string]map[string]struct{}
+ nodeCache map[string]nodeInfo
+ caCertPath string
+ caCertReloader *pkgtls.Reloader
+ prefix string
+ retryPolicy string
+ selfNode string
+ selfRole string
+ selfTier string
+ allowedRoles []databasev1.Role
+ healthCheckInterval time.Duration
+ writableProbeMu sync.Mutex
+ nodeCacheMu sync.RWMutex
+ tlsEnabled bool
}
// nodeInfo caches the resolved role and tier for a remote node.
@@ -227,6 +234,8 @@ func (p *pub) FlagSet() *run.FlagSet {
fs := run.NewFlagSet("queue-client")
fs.BoolVar(&p.tlsEnabled, prefixFlag("client-tls"), false,
fmt.Sprintf("enable client TLS for %s", p.prefix))
fs.StringVar(&p.caCertPath, prefixFlag("client-ca-cert"), "",
fmt.Sprintf("CA certificate file to verify the %s server", p.prefix))
+ fs.DurationVar(&p.healthCheckInterval,
prefixFlag("client-health-check-interval"), defaultHealthCheckInterval,
+ fmt.Sprintf("how often to re-check active %s nodes and evict
the unreachable ones; 0 disables the periodic check", p.prefix))
return fs
}
@@ -489,6 +498,10 @@ func NewWithoutMetadata(omr observability.MetricsRegistry)
queue.Client {
p := New(nil, databasev1.Role_ROLE_DATA)
pp := p.(*pub)
pp.log = logger.GetLogger("queue-client")
+ // No HealthCheckInterval here, unlike PreRun: this constructor skips
FlagSet, and its
+ // callers (lifecycle migration, tests) are short-lived and do not all
reach GracefulStop.
+ // Starting a background prober they never stop would leak a goroutine,
and the reactive
+ // failover path already covers a batch job's much narrower exposure.
pp.connMgr =
grpchelper.NewConnManager(grpchelper.ConnManagerConfig[*client]{
Handler: pp,
Logger: pp.log,
@@ -526,10 +539,11 @@ func (p *pub) PreRun(context.Context) error {
// Initialize connection manager with the pub as the handler
p.connMgr =
grpchelper.NewConnManager(grpchelper.ConnManagerConfig[*client]{
//nolint:contextcheck // health check runs in background goroutine
- Handler: p,
- Logger: p.log,
- RetryPolicy: p.retryPolicy,
- MaxRecvMsgSize: maxReceiveMessageSize,
+ Handler: p,
+ Logger: p.log,
+ RetryPolicy: p.retryPolicy,
+ MaxRecvMsgSize: maxReceiveMessageSize,
+ HealthCheckInterval: p.healthCheckInterval,
})
// Initialize CA certificate reloader if TLS is enabled and CA cert
path is provided
diff --git a/docs/operation/configuration.md b/docs/operation/configuration.md
index a4d193329..3675cb649 100644
--- a/docs/operation/configuration.md
+++ b/docs/operation/configuration.md
@@ -97,7 +97,13 @@ The following flags tune the BydbQL prepared-statement cache
on the query path.
These flags surface the queries behind cache misses and slow responses without
exposing high-cardinality query text as metric labels: Prometheus gets only two
counters (`bydbql_prepared_cache_total{result="miss"}` and
`bydbql_slow_query_total`), while the specific hot queries are logged.
- `--bydbql-slow-query-threshold duration`: End-to-end latency above which a
BydbQL query is counted as slow (increments `bydbql_slow_query_total`) and
tracked in the slow-query top-K; `0` disables slow-query tracking (default:
`1s`).
-- `--bydbql-topk-log-interval duration`: How often to log the hottest
cache-miss and slow queries. Counts are cumulative since process start. The
cache-miss list only shows templates re-parsed at least twice (`count>=2`):
every template misses once on its cold-start lookup, so a `count==1` entry is
benign and is filtered out — only repeatedly evicted-and-re-parsed (thrashing)
templates are surfaced. The slow-query list is ranked by peak latency
(`max_latency`) so a rarely-but-catastrophi [...]
+- `--bydbql-topk-log-interval duration`: How often to log the hottest
cache-miss and slow queries. The cache-miss list holds only templates that were
compiled more than once — either evicted and compiled again (thrashing), or too
large for the byte bound and therefore re-compiled on every single request; a
template's unavoidable first-ever compile is excluded at the source, so every
entry is actionable. The slow-query list is ranked by peak latency
(`max_latency`) so a rarely-but-catastr [...]
+- `--bydbql-topk-slow-ttl duration`: Drop a slow-query top-K entry whose query
has not been slow again for this long; `0` keeps entries for the process
lifetime (default: `24h`).
+- `--bydbql-topk-reparse-ttl duration`: Drop a cache-miss top-K entry whose
template has not been re-parsed again for this long; `0` keeps entries for the
process lifetime (default: `24h`).
+
+Each logged entry carries `last_seen`, and slow entries additionally carry
`max_latency_at`. Both trackers accumulate, so `max_latency` is a running peak
that can long outlive the incident that produced it: without these timestamps a
one-off startup spike keeps being reported as if it were current, and a reader
cannot tell a live problem from a stale one. The TTLs bound how long that can
happen — an entry whose query stops recurring disappears — while
`max_latency_at` dates the peak itse [...]
+
+Keep each TTL comfortably above `--bydbql-topk-log-interval`. An entry is only
ever reported by a dump, so a TTL shorter than the interval lets an entry
expire in the gap between two dumps and never be logged at all — the tracker
would go quiet not because nothing is wrong but because it forgets faster than
it reports. The defaults (`24h` against `5m`) leave a wide margin.
#### Diagnosing an ineffective BydbQL cache
@@ -238,6 +244,14 @@ These flags configure how every node talks to the schema
server.
- `--schema-property-client-tls`: Enable TLS for property schema client
connections.
- `--schema-property-client-ca-cert string`: CA certificate file to verify the
property schema server.
+### Queue client (liaison to data nodes)
+
+These flags configure the connections a liaison holds to the data nodes it
queries and writes to. `<prefix>` is `data` or `liaison` depending on which
peer set the client serves.
+
+- `--<prefix>-client-health-check-interval duration`: How often to re-check
the nodes the client currently considers active, evicting any that no longer
answer; `0` disables the periodic check (default: `10s`). The active set is
otherwise only validated when a node is admitted and, after that, when a
request happens to fail on it — so a node that dies stays routable until some
query picks it and pays that query's full timeout to discover it. Keep this
well below the distributed query tim [...]
+- `--<prefix>-client-tls`: Enable client TLS.
+- `--<prefix>-client-ca-cert string`: CA certificate file to verify the server.
+
### Other
- `-n, --name string`: Name of this service.