txn_offset_commit_request.go 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162
  1. package sarama
  2. type TxnOffsetCommitRequest struct {
  3. Version int16
  4. TransactionalID string
  5. GroupID string
  6. ProducerID int64
  7. ProducerEpoch int16
  8. Topics map[string][]*PartitionOffsetMetadata
  9. }
  10. func (t *TxnOffsetCommitRequest) encode(pe packetEncoder) error {
  11. if err := pe.putString(t.TransactionalID); err != nil {
  12. return err
  13. }
  14. if err := pe.putString(t.GroupID); err != nil {
  15. return err
  16. }
  17. pe.putInt64(t.ProducerID)
  18. pe.putInt16(t.ProducerEpoch)
  19. if err := pe.putArrayLength(len(t.Topics)); err != nil {
  20. return err
  21. }
  22. for topic, partitions := range t.Topics {
  23. if err := pe.putString(topic); err != nil {
  24. return err
  25. }
  26. if err := pe.putArrayLength(len(partitions)); err != nil {
  27. return err
  28. }
  29. for _, partition := range partitions {
  30. if err := partition.encode(pe, t.Version); err != nil {
  31. return err
  32. }
  33. }
  34. }
  35. return nil
  36. }
  37. func (t *TxnOffsetCommitRequest) decode(pd packetDecoder, version int16) (err error) {
  38. t.Version = version
  39. if t.TransactionalID, err = pd.getString(); err != nil {
  40. return err
  41. }
  42. if t.GroupID, err = pd.getString(); err != nil {
  43. return err
  44. }
  45. if t.ProducerID, err = pd.getInt64(); err != nil {
  46. return err
  47. }
  48. if t.ProducerEpoch, err = pd.getInt16(); err != nil {
  49. return err
  50. }
  51. n, err := pd.getArrayLength()
  52. if err != nil {
  53. return err
  54. }
  55. t.Topics = make(map[string][]*PartitionOffsetMetadata)
  56. for i := 0; i < n; i++ {
  57. topic, err := pd.getString()
  58. if err != nil {
  59. return err
  60. }
  61. m, err := pd.getArrayLength()
  62. if err != nil {
  63. return err
  64. }
  65. t.Topics[topic] = make([]*PartitionOffsetMetadata, m)
  66. for j := 0; j < m; j++ {
  67. partitionOffsetMetadata := new(PartitionOffsetMetadata)
  68. if err := partitionOffsetMetadata.decode(pd, version); err != nil {
  69. return err
  70. }
  71. t.Topics[topic][j] = partitionOffsetMetadata
  72. }
  73. }
  74. return nil
  75. }
  76. func (a *TxnOffsetCommitRequest) key() int16 {
  77. return 28
  78. }
  79. func (a *TxnOffsetCommitRequest) version() int16 {
  80. return a.Version
  81. }
  82. func (a *TxnOffsetCommitRequest) headerVersion() int16 {
  83. return 1
  84. }
  85. func (a *TxnOffsetCommitRequest) isValidVersion() bool {
  86. return a.Version >= 0 && a.Version <= 2
  87. }
  88. func (a *TxnOffsetCommitRequest) requiredVersion() KafkaVersion {
  89. switch a.Version {
  90. case 2:
  91. return V2_1_0_0
  92. case 1:
  93. return V2_0_0_0
  94. case 0:
  95. return V0_11_0_0
  96. default:
  97. return V2_1_0_0
  98. }
  99. }
  100. type PartitionOffsetMetadata struct {
  101. // Partition contains the index of the partition within the topic.
  102. Partition int32
  103. // Offset contains the message offset to be committed.
  104. Offset int64
  105. // LeaderEpoch contains the leader epoch of the last consumed record.
  106. LeaderEpoch int32
  107. // Metadata contains any associated metadata the client wants to keep.
  108. Metadata *string
  109. }
  110. func (p *PartitionOffsetMetadata) encode(pe packetEncoder, version int16) error {
  111. pe.putInt32(p.Partition)
  112. pe.putInt64(p.Offset)
  113. if version >= 2 {
  114. pe.putInt32(p.LeaderEpoch)
  115. }
  116. if err := pe.putNullableString(p.Metadata); err != nil {
  117. return err
  118. }
  119. return nil
  120. }
  121. func (p *PartitionOffsetMetadata) decode(pd packetDecoder, version int16) (err error) {
  122. if p.Partition, err = pd.getInt32(); err != nil {
  123. return err
  124. }
  125. if p.Offset, err = pd.getInt64(); err != nil {
  126. return err
  127. }
  128. if version >= 2 {
  129. if p.LeaderEpoch, err = pd.getInt32(); err != nil {
  130. return err
  131. }
  132. }
  133. if p.Metadata, err = pd.getNullableString(); err != nil {
  134. return err
  135. }
  136. return nil
  137. }