taos.go 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145
  1. // +build plugins
  2. package main
  3. import (
  4. "database/sql"
  5. "encoding/json"
  6. "fmt"
  7. "github.com/emqx/kuiper/common"
  8. "github.com/emqx/kuiper/xstream/api"
  9. _ "github.com/taosdata/driver-go/taosSql"
  10. "reflect"
  11. "strings"
  12. )
  13. type (
  14. taosConfig struct {
  15. Port int `json:"port"`
  16. Ip string `json:"ip"`
  17. User string `json:"user"`
  18. Password string `json:"password"`
  19. Database string `json:"database"`
  20. Table string `json:"table"`
  21. Fields []string `json:"fields"`
  22. }
  23. taosSink struct {
  24. conf *taosConfig
  25. db *sql.DB
  26. }
  27. )
  28. func (this *taosConfig) buildSql(ctx api.StreamContext, mapData map[string]interface{}) string {
  29. if 0 == len(mapData) {
  30. return ""
  31. }
  32. logger := ctx.GetLogger()
  33. var keys, vals []string
  34. for _, k := range this.Fields {
  35. if v, ok := mapData[k]; ok {
  36. keys = append(keys, k)
  37. if reflect.String == reflect.TypeOf(v).Kind() {
  38. vals = append(vals, fmt.Sprintf(`"%v"`, v))
  39. } else {
  40. vals = append(vals, fmt.Sprintf(`%v`, v))
  41. }
  42. } else {
  43. logger.Debug("not found field:", k)
  44. }
  45. }
  46. if 0 != len(keys) {
  47. if len(this.Fields) < len(mapData) {
  48. logger.Warnln("some of values will be ignored.")
  49. }
  50. return fmt.Sprintf(`INSERT INTO %s (%s)VALUES(%s);`, this.Table, strings.Join(keys, `,`), strings.Join(vals, `,`))
  51. }
  52. for k, v := range mapData {
  53. keys = append(keys, k)
  54. if reflect.String == reflect.TypeOf(v).Kind() {
  55. vals = append(vals, fmt.Sprintf(`"%v"`, v))
  56. } else {
  57. vals = append(vals, fmt.Sprintf(`%v`, v))
  58. }
  59. }
  60. if 0 != len(keys) {
  61. return fmt.Sprintf(`INSERT INTO %s (%s)VALUES(%s);`, this.Table, strings.Join(keys, `,`), strings.Join(vals, `,`))
  62. }
  63. return ""
  64. }
  65. func (m *taosSink) Configure(props map[string]interface{}) error {
  66. cfg := &taosConfig{}
  67. err := common.MapToStruct(props, cfg)
  68. if err != nil {
  69. return fmt.Errorf("read properties %v fail with error: %v", props, err)
  70. }
  71. if cfg.Ip == "" {
  72. cfg.Ip = "127.0.0.1"
  73. common.Log.Infof("Not find IP conf, will use default value '127.0.0.1'.")
  74. }
  75. if cfg.User == "" {
  76. cfg.User = "root"
  77. common.Log.Infof("Not find user conf, will use default value 'root'.")
  78. }
  79. if cfg.Password == "" {
  80. cfg.Password = "taosdata"
  81. common.Log.Infof("Not find password conf, will use default value 'taosdata'.")
  82. }
  83. if cfg.Database == "" {
  84. return fmt.Errorf("property database is required")
  85. }
  86. if cfg.Table == "" {
  87. return fmt.Errorf("property table is required")
  88. }
  89. m.conf = cfg
  90. return nil
  91. }
  92. func (m *taosSink) Open(ctx api.StreamContext) (err error) {
  93. logger := ctx.GetLogger()
  94. logger.Debug("Opening taos sink")
  95. url := fmt.Sprintf(`%s:%s@tcp(%s:%d)/%s`, m.conf.User, m.conf.Password, m.conf.Ip, m.conf.Port, m.conf.Database)
  96. m.db, err = sql.Open("taosSql", url)
  97. return err
  98. }
  99. func (m *taosSink) Collect(ctx api.StreamContext, item interface{}) error {
  100. logger := ctx.GetLogger()
  101. data, ok := item.([]byte)
  102. if !ok {
  103. logger.Debug("taos sink receive non string data")
  104. return nil
  105. }
  106. logger.Debugf("taos sink receive %s", item)
  107. var sliData []map[string]interface{}
  108. err := json.Unmarshal(data, &sliData)
  109. if nil != err {
  110. return err
  111. }
  112. for _, mapData := range sliData {
  113. sql := m.conf.buildSql(ctx, mapData)
  114. if 0 == len(sql) {
  115. continue
  116. }
  117. logger.Debugf(sql)
  118. rows, err := m.db.Query(sql)
  119. if err != nil {
  120. return err
  121. }
  122. rows.Close()
  123. }
  124. return nil
  125. }
  126. func (m *taosSink) Close(ctx api.StreamContext) error {
  127. if m.db != nil {
  128. return m.db.Close()
  129. }
  130. return nil
  131. }
  132. func Taos() api.Sink {
  133. return &taosSink{}
  134. }