Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
62 changes: 49 additions & 13 deletions pkg/wgo/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -537,12 +537,20 @@ func (c *WarpstreamClient) waitRefreshCooldown(delay, elapsed time.Duration) {
type rejectedTopicPartitionRecords struct {
topicPartitionRecords
err error
// miss is why the lookup found no agent, for the routing-miss counter.
miss routingMissReason
}

// routedGroup is a routed partition group plus how its agent was chosen.
type routedGroup struct {
promised[routedTopicPartitionRecords]
outcome routeOutcome
}

// routeRecords routes each partition once. A partition with no candidate is
// returned unsent. The first miss requests a metadata refresh.
func (c *WarpstreamClient) routeRecords(records []*kgo.Record, doneFor func(groupRecords []*kgo.Record) func(ProduceResult)) ([]promised[routedTopicPartitionRecords], []rejectedTopicPartitionRecords) {
groups := make(map[topicPartition]*promised[routedTopicPartitionRecords])
groups := make(map[topicPartition]*routedGroup)
order := make([]topicPartition, 0)
var rejectedByKey map[topicPartition]int
var rejected []rejectedTopicPartitionRecords
Expand All @@ -560,7 +568,7 @@ func (c *WarpstreamClient) routeRecords(records []*kgo.Record, doneFor func(grou
}
}

cands := c.demoter.Candidates(r.Topic, r.Partition, 1)
cands, outcome := c.demoter.candidatesWithRoute(r.Topic, r.Partition, 1)
if len(cands) == 0 {
if rejectedByKey == nil {
rejectedByKey = make(map[topicPartition]int)
Expand All @@ -573,30 +581,54 @@ func (c *WarpstreamClient) routeRecords(records []*kgo.Record, doneFor func(grou
partition: r.Partition,
records: []*kgo.Record{r},
},
err: fmt.Errorf("no agent assigned for topic %q partition %d", r.Topic, r.Partition),
err: fmt.Errorf("no agent assigned for topic %q partition %d", r.Topic, r.Partition),
miss: outcome.missReason(),
})
continue
}

groups[key] = &promised[routedTopicPartitionRecords]{
item: routedTopicPartitionRecords{
topicPartitionRecords: topicPartitionRecords{
topic: r.Topic,
partition: r.Partition,
records: []*kgo.Record{r},
groups[key] = &routedGroup{
promised: promised[routedTopicPartitionRecords]{
item: routedTopicPartitionRecords{
topicPartitionRecords: topicPartitionRecords{
topic: r.Topic,
partition: r.Partition,
records: []*kgo.Record{r},
},
nodeID: cands[0].NodeID,
nodeState: cands[0].State,
},
nodeID: cands[0].NodeID,
nodeState: cands[0].State,
},
outcome: outcome,
}
order = append(order, key)
}

var (
routes [routeSourceCount]int
misses [routingMissCount]int
)
out := make([]promised[routedTopicPartitionRecords], 0, len(order))
for _, key := range order {
g := groups[key]
g.done = doneFor(g.item.records)
out = append(out, *g)
out = append(out, g.promised)
if src, ok := g.outcome.routeSource(); ok {
routes[src] += len(g.item.records)
}
}
for i := range rejected {
misses[rejected[i].miss] += len(rejected[i].records)
}
for src, n := range routes {
if n > 0 {
c.metrics.partitionRoutes[src].Add(float64(n))
}
}
for reason, n := range misses {
if n > 0 {
c.metrics.routingMisses[reason].Add(float64(n))
}
}
return out, rejected
}
Expand All @@ -605,11 +637,15 @@ func (c *WarpstreamClient) routeRecords(records []*kgo.Record, doneFor func(grou
// skips the per-partition map and avoids heap-allocating an intermediate
// group. Returns an error if the record's partition has no known candidate.
func (c *WarpstreamClient) routeRecord(record *kgo.Record, done func(ProduceResult)) (promised[routedTopicPartitionRecords], error) {
cands := c.demoter.Candidates(record.Topic, record.Partition, 1)
cands, outcome := c.demoter.candidatesWithRoute(record.Topic, record.Partition, 1)
if len(cands) == 0 {
c.metrics.routingMisses[outcome.missReason()].Inc()
c.triggerRefresh()
return promised[routedTopicPartitionRecords]{}, fmt.Errorf("no agent assigned for topic %q partition %d", record.Topic, record.Partition)
}
if src, ok := outcome.routeSource(); ok {
c.metrics.partitionRoutes[src].Inc()
}
return promised[routedTopicPartitionRecords]{
item: routedTopicPartitionRecords{
topicPartitionRecords: topicPartitionRecords{
Expand Down
66 changes: 66 additions & 0 deletions pkg/wgo/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2466,3 +2466,69 @@ func (b *lockedBuffer) String() string {
defer b.mu.Unlock()
return b.buf.String()
}

func TestWarpstreamClient_RoutingMetricsReconcile(t *testing.T) {
const topic = "test-topic"
synctest.Test(t, func(t *testing.T) {
c, _, _, _ := newTestWarpstreamClient(t, topic, 4)
// The single fake broker is node 0. Partition 0 names it as leader,
// partition 1 has no entry (stand-in), partition 2 is unnamed and
// partitions past 3 do not exist.
live := c.pool.state.Load()
c.pool.state.Store(&poolState{
agents: live.agents,
topicIDs: live.topicIDs,
strategy: newDefaultPartitionAssignmentStrategy([]int32{0},
map[topicPartition]int32{{topic: topic, partition: 0}: 0},
nil,
map[topicPartition]struct{}{{topic: topic, partition: 2}: {}},
map[string]int32{topic: 4}),
})

rec := func(topic string, partition int32, size int) *kgo.Record {
return &kgo.Record{Topic: topic, Partition: partition, Value: make([]byte, size), Timestamp: time.Now()}
}
results := c.ProduceSync(t.Context(), []*kgo.Record{
rec(topic, 0, 1), rec(topic, 0, 1),
rec(topic, 1, 1), rec(topic, 1, 1), rec(topic, 1, 1),
rec(topic, 2, 1),
rec(topic, 9, 1), rec(topic, 9, 1),
rec("other", 0, 1),
rec(topic, 0, 2<<20), // larger than the test client's 1<<20 BatchMaxBytes
})
require.Len(t, results, 10)
for _, i := range []int{0, 1, 2, 3, 4} {
assert.NoError(t, results[i].Err)
}

var (
routes = testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceLeader]) + testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceStandIn])
misses float64
)
for _, m := range c.metrics.routingMisses {
misses += testutil.ToFloat64(m)
}
tooLarge := testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedRecordTooLarge))
noAgent := testutil.ToFloat64(c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned))

assert.Equal(t, float64(2), testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceLeader]))
assert.Equal(t, float64(3), testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceStandIn]))
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.routingMisses[routingMissNoLeader]))
assert.Equal(t, float64(2), testutil.ToFloat64(c.metrics.routingMisses[routingMissPartitionOutOfRange]))
assert.Equal(t, float64(1), testutil.ToFloat64(c.metrics.routingMisses[routingMissUnknownTopic]))
// Every submitted record is exactly one of: too large, routed, or a miss.
assert.Equal(t, testutil.ToFloat64(c.metrics.produceRecordsTotal), tooLarge+routes+misses)
assert.Equal(t, noAgent, misses)

// Produce counts one record per call. The rejected batch above nudged a
// refresh that may have replaced the injected strategy, so only the
// total is stable here.
var promised sync.WaitGroup
promised.Add(2)
c.Produce(t.Context(), rec(topic, 1, 1), func(*kgo.Record, error) { promised.Done() })
c.Produce(t.Context(), rec(topic, 9, 1), func(*kgo.Record, error) { promised.Done() })
promised.Wait()
assert.Equal(t, float64(6), testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceLeader])+testutil.ToFloat64(c.metrics.partitionRoutes[routeSourceStandIn]))
assert.Equal(t, float64(3), testutil.ToFloat64(c.metrics.routingMisses[routingMissPartitionOutOfRange]))
})
}
21 changes: 15 additions & 6 deletions pkg/wgo/demoter.go
Original file line number Diff line number Diff line change
Expand Up @@ -144,8 +144,16 @@ func NewDemoter(inner PartitionAssignmentStrategy, tracker AgentStatsReader, hea
// demoted and none is due for a probe, the natural primary is surfaced as a
// forced probe so we never refuse to route.
func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) []Agent {
agents, _ := d.candidatesWithRoute(topic, partition, maxCandidates)
return agents
}

// candidatesWithRoute is Candidates plus the inner strategy's classification of
// the lookup that produced the returned agents. Demotion only reorders or elides
// agents, so a demoted leader replaced here keeps that classification.
func (d *Demoter) candidatesWithRoute(topic string, partition int32, maxCandidates int) ([]Agent, routeOutcome) {
if maxCandidates <= 0 {
return nil
return nil, routeUnclassified
}
// Use the shared HealthCheckConfig parameters so the Demoter and the
// Hedger hit the same ClusterStats cache entry. cluster.SlowThreshold
Expand All @@ -155,7 +163,7 @@ func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) [
now := d.now()
clusterStats, hasClusterStats := d.tracker.ClusterStats(now, d.healthCfg.SlowMultiplier, d.healthCfg.FaultyThreshold)
if suppressed, _ := d.isDemotionSuppressed(clusterStats, hasClusterStats); suppressed {
return d.inner.Candidates(topic, partition, maxCandidates)
return candidatesOf(d.inner, topic, partition, maxCandidates)
}

// The returned list holds up to maxCandidates entries total
Expand All @@ -170,6 +178,7 @@ func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) [
const maxRetries = 6
var (
agents []Agent
route routeOutcome
noDemotedAgents bool
)
for retry, extra := 0, 2; retry < maxRetries; retry, extra = retry+1, extra*2 {
Expand All @@ -178,7 +187,7 @@ func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) [
nonDemoted = 0
)

agents = d.inner.Candidates(topic, partition, asked)
agents, route = candidatesOf(d.inner, topic, partition, asked)
noDemotedAgents = true

for _, c := range agents {
Expand All @@ -199,10 +208,10 @@ func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) [
}
}
if len(agents) == 0 {
return nil
return nil, route
}
if noDemotedAgents {
return agents[:min(len(agents), maxCandidates)]
return agents[:min(len(agents), maxCandidates)], route
}

candidates := make([]Agent, 0, maxCandidates)
Expand Down Expand Up @@ -239,7 +248,7 @@ func (d *Demoter) Candidates(topic string, partition int32, maxCandidates int) [
candidates = append(candidates, forced)
}

return candidates
return candidates, route

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Question: what is source meant to answer? Right now, when the Demoter skips a demoted leader and the record goes to a healthy alternate, it's still counted as source="leader", because route comes from the inner lookup. The help text says so and the "a demoted leader replaced by an alternate keeps the classification" test locks it in, so this looks deliberate.

The side effect is that during a demotion wave, leader + stand_in still looks like ~100% leader routing while most records actually go to alternates. If the counter is meant to show how the partition was resolved, this is fine. If it's also meant to show how often we route away from the named leader, it can't do that yet. That would take a third source (e.g. demoted_alternate), or a pointer to another metric that covers it. Which did you have in mind?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

source describes how the assignment strategy resolved the partition: leader means the named leader was available in the snapshot; stand_in means its leader entry was missing and the strategy picked a live agent.

The Demoter applies health policy to those candidates and may select an alternate. We preserve the original classification so this counter continues to measure metadata stand-in usage. Each record is counted once at initial routing.

You're right that it doesn't show how many records the Demoter redirects. I'd address that in a follow-up PR with a separate selection="original|alternate" label, preserving the leader-versus-stand-in distinction. The existing demotion and probe metrics provide context but don't count redirected records. I reworded the PR table to clarify the current meaning.

}

// isDemoted reports whether agent nodeID currently meets the demotion criteria.
Expand Down
Loading
Loading