kv_store.go 3.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143
  1. package states
  2. import (
  3. "encoding/gob"
  4. "fmt"
  5. "github.com/emqx/kuiper/common"
  6. "github.com/emqx/kuiper/common/kv"
  7. "github.com/emqx/kuiper/xstream/checkpoints"
  8. "path"
  9. "sync"
  10. )
  11. func init() {
  12. gob.Register(map[string]interface{}{})
  13. gob.Register(checkpoints.BufferOrEvent{})
  14. }
  15. //The manager for checkpoint storage.
  16. /**
  17. *** mapStore keys
  18. *** { "checkpoint1", "checkpoint2" ... "checkpointn" : The complete or incomplete snapshot
  19. */
  20. type KVStore struct {
  21. db kv.KeyValue
  22. mapStore *sync.Map //The current root store of a rule
  23. checkpoints []int64
  24. max int
  25. }
  26. //Store in path ./data/checkpoint/$ruleId
  27. //Store 2 things:
  28. //"checkpoints":A queue for completed checkpoint id
  29. //"$checkpointId":A map with key of checkpoint id and value of snapshot(gob serialized)
  30. //Assume each operator only has one instance
  31. func getKVStore(ruleId string) (*KVStore, error) {
  32. dr, _ := common.GetDataLoc()
  33. db := kv.GetDefaultKVStore(path.Join(dr, ruleId, "checkpoints"))
  34. s := &KVStore{db: db, max: 3, mapStore: &sync.Map{}}
  35. //read data from badger db
  36. if err := s.restore(); err != nil {
  37. return nil, err
  38. }
  39. return s, nil
  40. }
  41. func (s *KVStore) restore() error {
  42. err := s.db.Open()
  43. if err != nil {
  44. return err
  45. }
  46. defer s.db.Close()
  47. var cs []int64
  48. if ok, _ := s.db.Get(CheckpointListKey, &cs); ok {
  49. s.checkpoints = cs
  50. for _, c := range cs {
  51. var m map[string]interface{}
  52. if ok, _ := s.db.Get(fmt.Sprintf("%d", c), &m); ok {
  53. s.mapStore.Store(c, common.MapToSyncMap(m))
  54. } else {
  55. return fmt.Errorf("invalid checkpoint data: %v", c)
  56. }
  57. }
  58. }
  59. return nil
  60. }
  61. func (s *KVStore) SaveState(checkpointId int64, opId string, state map[string]interface{}) error {
  62. logger := common.Log
  63. logger.Debugf("Save state for checkpoint %d, op %s, value %v", checkpointId, opId, state)
  64. var cstore *sync.Map
  65. if v, ok := s.mapStore.Load(checkpointId); !ok {
  66. cstore = &sync.Map{}
  67. s.mapStore.Store(checkpointId, cstore)
  68. } else {
  69. if cstore, ok = v.(*sync.Map); !ok {
  70. return fmt.Errorf("invalid KVStore for checkpointId %d with value %v: should be *sync.Map type", checkpointId, v)
  71. }
  72. }
  73. cstore.Store(opId, state)
  74. return nil
  75. }
  76. func (s *KVStore) SaveCheckpoint(checkpointId int64) error {
  77. if v, ok := s.mapStore.Load(checkpointId); !ok {
  78. return fmt.Errorf("store for checkpoint %d not found", checkpointId)
  79. } else {
  80. if m, ok := v.(*sync.Map); !ok {
  81. return fmt.Errorf("invalid KVStore for checkpointId %d with value %v: should be *sync.Map type", checkpointId, v)
  82. } else {
  83. err := s.db.Open()
  84. if err != nil {
  85. return fmt.Errorf("save checkpoint err: %v", err)
  86. }
  87. defer s.db.Close()
  88. err = s.db.Set(fmt.Sprintf("%d", checkpointId), common.SyncMapToMap(m))
  89. if err != nil {
  90. return fmt.Errorf("save checkpoint err: %v", err)
  91. }
  92. m.Delete(checkpointId)
  93. s.checkpoints = append(s.checkpoints, checkpointId)
  94. //TODO is the order promised?
  95. if len(s.checkpoints) > s.max {
  96. cp := s.checkpoints[0]
  97. s.checkpoints = s.checkpoints[1:]
  98. go func() {
  99. _ = s.db.Delete(fmt.Sprintf("%d", cp))
  100. }()
  101. }
  102. err = s.db.Set(CheckpointListKey, s.checkpoints)
  103. if err != nil {
  104. return fmt.Errorf("save checkpoint err: %v", err)
  105. }
  106. }
  107. }
  108. return nil
  109. }
  110. //Only run in the initialization
  111. func (s *KVStore) GetOpState(opId string) (*sync.Map, error) {
  112. if len(s.checkpoints) > 0 {
  113. if v, ok := s.mapStore.Load(s.checkpoints[len(s.checkpoints)-1]); ok {
  114. if cstore, ok := v.(*sync.Map); !ok {
  115. return nil, fmt.Errorf("invalid state %v stored for op %s: data type is not *sync.Map", v, opId)
  116. } else {
  117. if sm, ok := cstore.Load(opId); ok {
  118. switch m := sm.(type) {
  119. case *sync.Map:
  120. return m, nil
  121. case map[string]interface{}:
  122. return common.MapToSyncMap(m), nil
  123. default:
  124. return nil, fmt.Errorf("invalid state %v stored for op %s: data type is not *sync.Map", sm, opId)
  125. }
  126. }
  127. }
  128. } else {
  129. return nil, fmt.Errorf("store for checkpoint %d not found", s.checkpoints[len(s.checkpoints)-1])
  130. }
  131. }
  132. return &sync.Map{}, nil
  133. }