Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,8 @@ For every record the client asks a `PartitionAssignmentStrategy` for an ordered

A partition's leader can briefly go missing from that map during a Metadata refresh, even though nothing is actually wrong. When the topic is known, the client picks a live agent for that partition instead of treating it as unroutable — any agent can serve any partition. That call succeeds without waiting. The refresh that excluded the leader asks for another fetch on its own, even when no Produce is in flight. Further exclusions double the wait from `OnDemandMetadataRefreshInterval` up to `MetadataRefreshInterval`. Repeated routing misses climb the same way. While on-demand refreshes keep being requested, the periodic refresh does not run, so the gap stays at the ceiling until they stop. A refresh that excludes nobody stops the follow-ups, and the next periodic refresh puts the wait back. A topic is known when another partition still has a leader, and also when this refresh listed the topic's partitions and kept none of the leaders it named. A partition Metadata reported with no leader (`Leader` below 0) is left unroutable, including when a sibling partition still has one, so that produce fails and still gets an on-demand refresh. A topic that has never appeared in Metadata, or that came back with an error, is left alone for the same reason.

`Produce` rejects only the record whose partition has no candidate. `ProduceSync` does the same within one call: a partition with a candidate is buffered and the call waits for it, and a partition with no candidate is rejected on its own and is not sent. That rejection does not fail the other records.

### Buffering: linger by destination agent, not by partition

Records are buffered through a `ClusterRecordBuffer`, which bins them by the destination agent picked at routing time, then through a per-agent `AgentRecordBuffer`, which applies a configurable linger window before flushing. Each flush ships one Produce request to one agent carrying batches for as many partitions as the buffer accumulated.
Expand Down
5 changes: 4 additions & 1 deletion docs/internal/tracing.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,10 @@ franz-go uses:
- **Unbuffered** fires before the caller observes a record's outcome. For `Produce` the
promise is wrapped so the hook runs just before it, mirroring franz-go (unbuffered hook,
then promise). For `ProduceSync` the hooks fire, in input order, on the calling goroutine —
after it has finalized every result.
after it has finalized every result. A call that rejects some records for routing and
produces the rest still fires unbuffered once per record, at that return. Spans for the
rejected records stay open until the accepted records finish. Rejected records do not
get an earlier unbuffered hook.

Callers pass their tracer via the existing `WithHooks`, exactly as for a franz-go client,
so this needs no new API and no OpenTelemetry dependency in the client — it relies only on
Expand Down
121 changes: 83 additions & 38 deletions pkg/wgo/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -290,40 +290,57 @@ func (c *WarpstreamClient) ProduceSync(ctx context.Context, records []*kgo.Recor
}

wg.Add(len(okRecords))
// Last input index for each pointer. Completions write that slot only.
indexOf := make(map[*kgo.Record]int, len(okRecords))
for _, idx := range okIndices {
indexOf[records[idx]] = idx
}
// A repeated pointer shares that slot's outcome. Copy it to the other
// positions on return, before the unbuffered hooks above.
if len(indexOf) < len(okRecords) {
defer func() {
for _, idx := range okIndices {
canon := indexOf[records[idx]]
if idx != canon {
results[idx] = results[canon]
}
}
}()
}
write := func(recs []*kgo.Record, err error) {
for _, r := range recs {
results[indexOf[r]] = kgo.ProduceResult{Record: r, Err: err}
wg.Done()
}
}

routed, err := c.routeRecords(okRecords, func(groupRecords []*kgo.Record) func(ProduceResult) {
routed, rejected := c.routeRecords(okRecords, func(groupRecords []*kgo.Record) func(ProduceResult) {
return perPartitionDone(groupRecords[0].Topic, groupRecords[0].Partition, groupRecords, func(err error) {
if err != nil {
// Post-dispatch failure, resolved uniformly for the whole
// partition group; pre-dispatch rejections never reach here.
c.metrics.produceRecordsFailedTotal.Add(float64(len(groupRecords)))
}
for _, r := range groupRecords {
results[indexOf[r]] = kgo.ProduceResult{Record: r, Err: err}
wg.Done()
}
write(groupRecords, err)
})
})
if err != nil {
// One record had no known candidate. Fail the whole batch
// uniformly: every ok record gets the same error.
c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned).Add(float64(len(okIndices)))
for _, i := range okIndices {
results[i] = kgo.ProduceResult{Record: records[i], Err: err}

if len(rejected) > 0 {
for _, rg := range rejected {
c.metrics.produceRecordsRejectedTotal.WithLabelValues(produceRejectedNoAgentAssigned).Add(float64(len(rg.records)))
write(rg.records, rg.err)
}
}
if len(routed) == 0 {
return results
}

// Stamp each record's produce time only after routing succeeds, so a failed
// produce leaves the caller's records unchanged. A single now keeps records
// buffered together on one produce timestamp. Mirrors franz-go's bufferRecord.
// Stamp only accepted records with unset timestamps, using one shared now.
now := time.Now()
for _, r := range okRecords {
ensureRecordTimestamp(r, now)
for _, g := range routed {
for _, r := range g.item.records {
ensureRecordTimestamp(r, now)
}
}

c.buffer.MultiAdd(ctx, routed)
Expand Down Expand Up @@ -540,43 +557,71 @@ func (c *WarpstreamClient) waitRefreshCooldown(delay, elapsed time.Duration) {
}
}

// routeRecords groups records by (topic, partition), stamps each group with
// its initial destination NodeID and mints the per-group done callback.
// Returns an error if any record's partition has no known candidate.
func (c *WarpstreamClient) routeRecords(records []*kgo.Record, doneFor func(groupRecords []*kgo.Record) func(ProduceResult)) ([]promised[routedTopicPartitionRecords], error) {
type rejectedTopicPartitionRecords struct {
topicPartitionRecords
err error
}

// 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])
order := make([]topicPartition, 0)
var rejectedByKey map[topicPartition]int
var rejected []rejectedTopicPartitionRecords

for _, r := range records {
key := topicPartition{topic: r.Topic, partition: r.Partition}
g, ok := groups[key]
if !ok {
cands := c.demoter.Candidates(r.Topic, r.Partition, 1)
if len(cands) == 0 {
if g, ok := groups[key]; ok {
g.item.records = append(g.item.records, r)
continue
}
if rejectedByKey != nil {
if i, ok := rejectedByKey[key]; ok {
rejected[i].records = append(rejected[i].records, r)
continue
}
}

cands := c.demoter.Candidates(r.Topic, r.Partition, 1)
if len(cands) == 0 {
if rejectedByKey == nil {
rejectedByKey = make(map[topicPartition]int)
c.triggerRefresh()
return nil, fmt.Errorf("no agent assigned for topic %q partition %d", r.Topic, r.Partition)
}
g = &promised[routedTopicPartitionRecords]{
item: routedTopicPartitionRecords{
topicPartitionRecords: topicPartitionRecords{
topic: r.Topic,
partition: r.Partition,
},
nodeID: cands[0].NodeID,
nodeState: cands[0].State,
rejectedByKey[key] = len(rejected)
rejected = append(rejected, rejectedTopicPartitionRecords{
topicPartitionRecords: topicPartitionRecords{
topic: r.Topic,
partition: r.Partition,
records: []*kgo.Record{r},
},
}
groups[key] = g
order = append(order, key)
err: fmt.Errorf("no agent assigned for topic %q partition %d", r.Topic, r.Partition),
})
continue
}

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

out := make([]promised[routedTopicPartitionRecords], 0, len(order))
for _, key := range order {
g := groups[key]
g.done = doneFor(g.item.records)
out = append(out, *g)
}
return out, nil
return out, rejected
}

// routeRecord is the single-record specialisation of routeRecords: it
Expand Down
Loading
Loading