random.go 3.4 KB

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