-
Notifications
You must be signed in to change notification settings - Fork 18
/
offset_fetch_request.go
125 lines (101 loc) · 2.73 KB
/
offset_fetch_request.go
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
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
package healer
/*
v0 and v1 are identical on the wire, but v0 (supported in 0.8.1 or later) reads offsets from zookeeper, while v1 (supported in 0.8.2 or later) reads offsets from kafka.
OffsetFetch Request (Version: 0) => group_id [topics]
group_id => STRING
topics => topic [partitions]
topic => STRING
partitions => partition
partition => INT32
group_id The unique group identifier
topics Topics to fetch offsets.
topic Name of topic
partitions Partitions to fetch offsets.
partition Topic partition id
*/
import (
"encoding/binary"
)
type OffsetFetchRequestTopic struct {
Topic string
Partitions []int32
}
type OffsetFetchRequest struct {
*RequestHeader
GroupID string
Topics []*OffsetFetchRequestTopic
}
// request only ONE topic
func NewOffsetFetchRequest(apiVersion uint16, clientID, groupID string) *OffsetFetchRequest {
requestHeader := &RequestHeader{
APIKey: API_OffsetFetchRequest,
APIVersion: apiVersion,
ClientID: clientID,
}
r := &OffsetFetchRequest{
RequestHeader: requestHeader,
GroupID: groupID,
}
r.Topics = make([]*OffsetFetchRequestTopic, 0)
return r
}
func (r *OffsetFetchRequest) AddPartiton(topic string, partitionID int32) {
if r.Topics == nil {
r.Topics = make([]*OffsetFetchRequestTopic, 0)
}
var theTopic *OffsetFetchRequestTopic = nil
for _, t := range r.Topics {
if t.Topic == topic {
theTopic = t
break
}
}
if theTopic == nil {
theTopic = &OffsetFetchRequestTopic{
Topic: topic,
}
r.Topics = append(r.Topics, theTopic)
}
for _, p := range theTopic.Partitions {
if p == partitionID {
return
}
}
theTopic.Partitions = append(theTopic.Partitions, partitionID)
return
}
func (r *OffsetFetchRequest) Length() int {
l := r.RequestHeader.length()
l += 2 + len(r.GroupID)
l += 4
for _, t := range r.Topics {
l += 2 + len(t.Topic)
l += 4 + 4*len(t.Partitions)
}
return l
}
func (r *OffsetFetchRequest) Encode(version uint16) []byte {
requestLength := r.Length()
payload := make([]byte, 4+requestLength)
offset := 0
binary.BigEndian.PutUint32(payload[offset:], uint32(requestLength))
offset += 4
offset += r.RequestHeader.Encode(payload[offset:])
binary.BigEndian.PutUint16(payload[offset:], uint16(len(r.GroupID)))
offset += 2
offset += copy(payload[offset:], r.GroupID)
binary.BigEndian.PutUint32(payload[offset:], uint32(len(r.Topics)))
offset += 4
for _, t := range r.Topics {
binary.BigEndian.PutUint16(payload[offset:], uint16(len(t.Topic)))
offset += 2
offset += copy(payload[offset:], t.Topic)
binary.BigEndian.PutUint32(payload[offset:], uint32(len(t.Partitions)))
offset += 4
for _, p := range t.Partitions {
binary.BigEndian.PutUint32(payload[offset:], uint32(p))
offset += 4
}
}
return payload
}