offset_fetch_request.go 1.5 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071
  1. package sarama
  2. type OffsetFetchRequest struct {
  3. ConsumerGroup string
  4. Version int16
  5. partitions map[string][]int32
  6. }
  7. func (r *OffsetFetchRequest) encode(pe packetEncoder) (err error) {
  8. if r.Version < 0 || r.Version > 1 {
  9. return PacketEncodingError{"invalid or unsupported OffsetFetchRequest version field"}
  10. }
  11. if err = pe.putString(r.ConsumerGroup); err != nil {
  12. return err
  13. }
  14. if err = pe.putArrayLength(len(r.partitions)); err != nil {
  15. return err
  16. }
  17. for topic, partitions := range r.partitions {
  18. if err = pe.putString(topic); err != nil {
  19. return err
  20. }
  21. if err = pe.putInt32Array(partitions); err != nil {
  22. return err
  23. }
  24. }
  25. return nil
  26. }
  27. func (r *OffsetFetchRequest) decode(pd packetDecoder) (err error) {
  28. if r.ConsumerGroup, err = pd.getString(); err != nil {
  29. return err
  30. }
  31. partitionCount, err := pd.getArrayLength()
  32. if err != nil {
  33. return err
  34. }
  35. if partitionCount == 0 {
  36. return nil
  37. }
  38. r.partitions = make(map[string][]int32)
  39. for i := 0; i < partitionCount; i++ {
  40. topic, err := pd.getString()
  41. if err != nil {
  42. return err
  43. }
  44. partitions, err := pd.getInt32Array()
  45. if err != nil {
  46. return err
  47. }
  48. r.partitions[topic] = partitions
  49. }
  50. return nil
  51. }
  52. func (r *OffsetFetchRequest) key() int16 {
  53. return 9
  54. }
  55. func (r *OffsetFetchRequest) version() int16 {
  56. return r.Version
  57. }
  58. func (r *OffsetFetchRequest) AddPartition(topic string, partitionID int32) {
  59. if r.partitions == nil {
  60. r.partitions = make(map[string][]int32)
  61. }
  62. r.partitions[topic] = append(r.partitions[topic], partitionID)
  63. }