You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

109 lines
3.5 KiB

  1. package segment
  2. import (
  3. flatbuffers "github.com/google/flatbuffers/go"
  4. "github.com/seaweedfs/seaweedfs/weed/pb/message_fbs"
  5. )
  6. type MessageBatchBuilder struct {
  7. b *flatbuffers.Builder
  8. producerId int32
  9. producerEpoch int32
  10. segmentId int32
  11. flags int32
  12. messageOffsets []flatbuffers.UOffsetT
  13. segmentSeqBase int64
  14. segmentSeqLast int64
  15. tsMsBase int64
  16. tsMsLast int64
  17. }
  18. func NewMessageBatchBuilder(b *flatbuffers.Builder,
  19. producerId int32,
  20. producerEpoch int32,
  21. segmentId int32,
  22. flags int32) *MessageBatchBuilder {
  23. b.Reset()
  24. return &MessageBatchBuilder{
  25. b: b,
  26. producerId: producerId,
  27. producerEpoch: producerEpoch,
  28. segmentId: segmentId,
  29. flags: flags,
  30. }
  31. }
  32. func (builder *MessageBatchBuilder) AddMessage(segmentSeq int64, tsMs int64, properties map[string][]byte, key []byte, value []byte) {
  33. if builder.segmentSeqBase == 0 {
  34. builder.segmentSeqBase = segmentSeq
  35. }
  36. builder.segmentSeqLast = segmentSeq
  37. if builder.tsMsBase == 0 {
  38. builder.tsMsBase = tsMs
  39. }
  40. builder.tsMsLast = tsMs
  41. var names, values, pairs []flatbuffers.UOffsetT
  42. for k, v := range properties {
  43. names = append(names, builder.b.CreateString(k))
  44. values = append(values, builder.b.CreateByteVector(v))
  45. }
  46. for i, _ := range names {
  47. message_fbs.NameValueStart(builder.b)
  48. message_fbs.NameValueAddName(builder.b, names[i])
  49. message_fbs.NameValueAddValue(builder.b, values[i])
  50. pair := message_fbs.NameValueEnd(builder.b)
  51. pairs = append(pairs, pair)
  52. }
  53. message_fbs.MessageStartPropertiesVector(builder.b, len(properties))
  54. for i := len(pairs) - 1; i >= 0; i-- {
  55. builder.b.PrependUOffsetT(pairs[i])
  56. }
  57. propOffset := builder.b.EndVector(len(properties))
  58. keyOffset := builder.b.CreateByteVector(key)
  59. valueOffset := builder.b.CreateByteVector(value)
  60. message_fbs.MessageStart(builder.b)
  61. message_fbs.MessageAddSeqDelta(builder.b, int32(segmentSeq-builder.segmentSeqBase))
  62. message_fbs.MessageAddTsMsDelta(builder.b, int32(tsMs-builder.tsMsBase))
  63. message_fbs.MessageAddProperties(builder.b, propOffset)
  64. message_fbs.MessageAddKey(builder.b, keyOffset)
  65. message_fbs.MessageAddData(builder.b, valueOffset)
  66. messageOffset := message_fbs.MessageEnd(builder.b)
  67. builder.messageOffsets = append(builder.messageOffsets, messageOffset)
  68. }
  69. func (builder *MessageBatchBuilder) BuildMessageBatch() {
  70. message_fbs.MessageBatchStartMessagesVector(builder.b, len(builder.messageOffsets))
  71. for i := len(builder.messageOffsets) - 1; i >= 0; i-- {
  72. builder.b.PrependUOffsetT(builder.messageOffsets[i])
  73. }
  74. messagesOffset := builder.b.EndVector(len(builder.messageOffsets))
  75. message_fbs.MessageBatchStart(builder.b)
  76. message_fbs.MessageBatchAddProducerId(builder.b, builder.producerId)
  77. message_fbs.MessageBatchAddProducerEpoch(builder.b, builder.producerEpoch)
  78. message_fbs.MessageBatchAddSegmentId(builder.b, builder.segmentId)
  79. message_fbs.MessageBatchAddFlags(builder.b, builder.flags)
  80. message_fbs.MessageBatchAddSegmentSeqBase(builder.b, builder.segmentSeqBase)
  81. message_fbs.MessageBatchAddSegmentSeqMaxDelta(builder.b, int32(builder.segmentSeqLast-builder.segmentSeqBase))
  82. message_fbs.MessageBatchAddTsMsBase(builder.b, builder.tsMsBase)
  83. message_fbs.MessageBatchAddTsMsMaxDelta(builder.b, int32(builder.tsMsLast-builder.tsMsBase))
  84. message_fbs.MessageBatchAddMessages(builder.b, messagesOffset)
  85. messageBatch := message_fbs.MessageBatchEnd(builder.b)
  86. builder.b.Finish(messageBatch)
  87. }
  88. func (builder *MessageBatchBuilder) GetBytes() []byte {
  89. return builder.b.FinishedBytes()
  90. }