random.go 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135
  1. package main
  2. import (
  3. "bytes"
  4. "encoding/json"
  5. "fmt"
  6. "github.com/emqx/kuiper/pkg/api"
  7. "github.com/emqx/kuiper/pkg/cast"
  8. "github.com/emqx/kuiper/pkg/message"
  9. "math/rand"
  10. "strings"
  11. "time"
  12. )
  13. const dedupStateKey = "input"
  14. type randomSourceConfig struct {
  15. Interval int `json:"interval"`
  16. Seed int `json:"seed"`
  17. Pattern map[string]interface{} `json:"pattern"`
  18. // how long will the source trace for deduplication. If 0, deduplicate is disabled; if negative, deduplicate will be the whole life time
  19. Deduplicate int `json:"deduplicate"`
  20. Format string `json:"format"`
  21. }
  22. //Emit data randomly with only a string field
  23. type randomSource struct {
  24. conf *randomSourceConfig
  25. list [][]byte
  26. }
  27. func (s *randomSource) Configure(topic string, props map[string]interface{}) error {
  28. cfg := &randomSourceConfig{}
  29. err := cast.MapToStruct(props, cfg)
  30. if err != nil {
  31. return fmt.Errorf("read properties %v fail with error: %v", props, err)
  32. }
  33. if cfg.Interval <= 0 {
  34. return fmt.Errorf("source `random` property `interval` must be a positive integer but got %d", cfg.Interval)
  35. }
  36. if cfg.Pattern == nil {
  37. return fmt.Errorf("source `random` property `pattern` is required")
  38. }
  39. if cfg.Interval <= 0 {
  40. return fmt.Errorf("source `random` property `seed` must be a positive integer but got %d", cfg.Seed)
  41. }
  42. if strings.ToLower(cfg.Format) != message.FormatJson {
  43. return fmt.Errorf("random source only supports `json` format")
  44. }
  45. s.conf = cfg
  46. return nil
  47. }
  48. func (s *randomSource) Open(ctx api.StreamContext, consumer chan<- api.SourceTuple, errCh chan<- error) {
  49. logger := ctx.GetLogger()
  50. logger.Debugf("open random source with deduplicate %d", s.conf.Deduplicate)
  51. if s.conf.Deduplicate != 0 {
  52. list, err := ctx.GetState(dedupStateKey)
  53. if err != nil {
  54. errCh <- err
  55. return
  56. }
  57. if list == nil {
  58. list = make([][]byte, 0)
  59. } else {
  60. if l, ok := list.([][]byte); ok {
  61. logger.Debugf("restore list %v", l)
  62. s.list = l
  63. } else {
  64. s.list = make([][]byte, 0)
  65. logger.Warnf("random source gets invalid state, ignore it")
  66. }
  67. }
  68. }
  69. t := time.NewTicker(time.Duration(s.conf.Interval) * time.Millisecond)
  70. defer t.Stop()
  71. for {
  72. select {
  73. case <-t.C:
  74. next := randomize(s.conf.Pattern, s.conf.Seed)
  75. if s.conf.Deduplicate != 0 && s.isDup(ctx, next) {
  76. logger.Debugf("find duplicate")
  77. continue
  78. }
  79. logger.Debugf("Send out data %v", next)
  80. consumer <- api.NewDefaultSourceTuple(next, nil)
  81. case <-ctx.Done():
  82. return
  83. }
  84. }
  85. }
  86. func randomize(p map[string]interface{}, seed int) map[string]interface{} {
  87. r := make(map[string]interface{})
  88. for k, v := range p {
  89. //TODO other data types
  90. vi, err := cast.ToInt(v, cast.STRICT)
  91. if err != nil {
  92. break
  93. }
  94. r[k] = vi + rand.Intn(seed)
  95. }
  96. return r
  97. }
  98. func (s *randomSource) isDup(ctx api.StreamContext, next map[string]interface{}) bool {
  99. logger := ctx.GetLogger()
  100. ns, err := json.Marshal(next)
  101. if err != nil {
  102. logger.Warnf("invalid input data %v", next)
  103. return true
  104. }
  105. for _, ps := range s.list {
  106. if bytes.Compare(ns, ps) == 0 {
  107. logger.Debugf("got duplicate %s", ns)
  108. return true
  109. }
  110. }
  111. logger.Debugf("no duplicate %s", ns)
  112. if s.conf.Deduplicate > 0 && len(s.list) >= s.conf.Deduplicate {
  113. s.list = s.list[1:]
  114. }
  115. s.list = append(s.list, ns)
  116. ctx.PutState(dedupStateKey, s.list)
  117. return false
  118. }
  119. func (s *randomSource) Close(_ api.StreamContext) error {
  120. return nil
  121. }
  122. func Random() api.Source {
  123. return &randomSource{}
  124. }