message_serde.go 3.3 KB

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