-
Notifications
You must be signed in to change notification settings - Fork 860
Expand file tree
/
Copy pathmetadata_refresh.go
More file actions
89 lines (84 loc) · 2.53 KB
/
Copy pathmetadata_refresh.go
File metadata and controls
89 lines (84 loc) · 2.53 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
package kafka
import (
"errors"
"github.com/segmentio/kafka-go/protocol"
fetchAPI "github.com/segmentio/kafka-go/protocol/fetch"
produceAPI "github.com/segmentio/kafka-go/protocol/produce"
)
// isStaleMetadataError reports whether a kafka error code indicates that the
// cached cluster metadata is likely out of date and should be refreshed.
func isStaleMetadataError(err Error) bool {
switch err {
case UnknownTopicOrPartition,
LeaderNotAvailable,
NotLeaderForPartition,
ReplicaNotAvailable,
BrokerNotAvailable,
KafkaStorageError,
FencedLeaderEpoch,
UnknownLeaderEpoch,
UnknownTopicID:
return true
default:
return false
}
}
// errorRequiresMetadataRefresh reports whether a transport-level error returned
// by a failed round trip likely indicates stale cluster metadata: transient
// network errors, routing failures over a stale layout (unknown topic, partition
// or leader), or kafka errors that signal the metadata is out of date.
func errorRequiresMetadataRefresh(err error) bool {
if isTransientNetworkError(err) {
return true
}
// Routing failures raised while resolving the target broker from the cached
// cluster layout indicate the layout is stale and should be refreshed.
if errors.Is(err, protocol.ErrNoTopic) ||
errors.Is(err, protocol.ErrNoPartition) ||
errors.Is(err, protocol.ErrNoLeader) {
return true
}
var kafkaErr Error
if errors.As(err, &kafkaErr) {
return isStaleMetadataError(kafkaErr)
}
return false
}
// responseRequiresMetadataRefresh inspects a successful response body and
// reports whether any partition (or, for Fetch, the session-level error code)
// reported a stale-metadata error. Only response types routed to specific
// partition leaders are inspected, since those are impacted by leader changes.
func responseRequiresMetadataRefresh(r Response) bool {
switch resp := r.(type) {
case *produceAPI.Response:
if resp == nil {
return false
}
for i := range resp.Topics {
partitions := resp.Topics[i].Partitions
for j := range partitions {
if isStaleMetadataError(Error(partitions[j].ErrorCode)) {
return true
}
}
}
case *fetchAPI.Response:
if resp == nil {
return false
}
// Fetch v7+ may report a session-level error in addition to the
// per-partition error codes.
if isStaleMetadataError(Error(resp.ErrorCode)) {
return true
}
for i := range resp.Topics {
partitions := resp.Topics[i].Partitions
for j := range partitions {
if isStaleMetadataError(Error(partitions[j].ErrorCode)) {
return true
}
}
}
}
return false
}