taos.go 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180
  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. ProvideTs bool `json:"provideTs"`
  16. Port int `json:"port"`
  17. Ip string `json:"ip"`
  18. User string `json:"user"`
  19. Password string `json:"password"`
  20. Database string `json:"database"`
  21. Table string `json:"table"`
  22. TsFieldName string `json:"tsFieldName"`
  23. Fields []string `json:"fields"`
  24. }
  25. taosSink struct {
  26. conf *taosConfig
  27. db *sql.DB
  28. }
  29. )
  30. func (this *taosConfig) delTsField() {
  31. var auxFields []string
  32. for _, v := range this.Fields {
  33. if v != this.TsFieldName {
  34. auxFields = append(auxFields, v)
  35. }
  36. }
  37. this.Fields = auxFields
  38. }
  39. func (this *taosConfig) buildSql(ctx api.StreamContext, mapData map[string]interface{}) (string, error) {
  40. if 0 == len(mapData) {
  41. return "", fmt.Errorf("data is empty.")
  42. }
  43. if 0 == len(this.TsFieldName) {
  44. return "", fmt.Errorf("tsFieldName is empty.")
  45. }
  46. logger := ctx.GetLogger()
  47. var keys, vals []string
  48. if this.ProvideTs {
  49. if v, ok := mapData[this.TsFieldName]; !ok {
  50. return "", fmt.Errorf("Timestamp field not found : %s.", this.TsFieldName)
  51. } else {
  52. keys = append(keys, this.TsFieldName)
  53. vals = append(vals, fmt.Sprintf(`"%v"`, v))
  54. delete(mapData, this.TsFieldName)
  55. this.delTsField()
  56. }
  57. } else {
  58. vals = append(vals, "now")
  59. keys = append(keys, this.TsFieldName)
  60. }
  61. for _, k := range this.Fields {
  62. if v, ok := mapData[k]; ok {
  63. keys = append(keys, k)
  64. if reflect.String == reflect.TypeOf(v).Kind() {
  65. vals = append(vals, fmt.Sprintf(`"%v"`, v))
  66. } else {
  67. vals = append(vals, fmt.Sprintf(`%v`, v))
  68. }
  69. } else {
  70. logger.Warnln("not found field:", k)
  71. }
  72. }
  73. if 0 != len(this.Fields) {
  74. if len(this.Fields) < len(mapData) {
  75. logger.Warnln("some of values will be ignored.")
  76. }
  77. return fmt.Sprintf(`INSERT INTO %s (%s)VALUES(%s);`, this.Table, strings.Join(keys, `,`), strings.Join(vals, `,`)), nil
  78. }
  79. for k, v := range mapData {
  80. keys = append(keys, k)
  81. if reflect.String == reflect.TypeOf(v).Kind() {
  82. vals = append(vals, fmt.Sprintf(`"%v"`, v))
  83. } else {
  84. vals = append(vals, fmt.Sprintf(`%v`, v))
  85. }
  86. }
  87. if 0 != len(keys) {
  88. return fmt.Sprintf(`INSERT INTO %s (%s)VALUES(%s);`, this.Table, strings.Join(keys, `,`), strings.Join(vals, `,`)), nil
  89. }
  90. return "", nil
  91. }
  92. func (m *taosSink) Configure(props map[string]interface{}) error {
  93. cfg := &taosConfig{}
  94. err := common.MapToStruct(props, cfg)
  95. if err != nil {
  96. return fmt.Errorf("read properties %v fail with error: %v", props, err)
  97. }
  98. if cfg.Ip == "" {
  99. cfg.Ip = "127.0.0.1"
  100. common.Log.Infof("Not find IP conf, will use default value '127.0.0.1'.")
  101. }
  102. if cfg.User == "" {
  103. cfg.User = "root"
  104. common.Log.Infof("Not find user conf, will use default value 'root'.")
  105. }
  106. if cfg.Password == "" {
  107. cfg.Password = "taosdata"
  108. common.Log.Infof("Not find password conf, will use default value 'taosdata'.")
  109. }
  110. if cfg.Database == "" {
  111. return fmt.Errorf("property database is required")
  112. }
  113. if cfg.Table == "" {
  114. return fmt.Errorf("property table is required")
  115. }
  116. if cfg.TsFieldName == "" {
  117. return fmt.Errorf("property TsFieldName is required")
  118. }
  119. m.conf = cfg
  120. return nil
  121. }
  122. func (m *taosSink) Open(ctx api.StreamContext) (err error) {
  123. logger := ctx.GetLogger()
  124. logger.Debug("Opening taos sink")
  125. url := fmt.Sprintf(`%s:%s@tcp(%s:%d)/%s`, m.conf.User, m.conf.Password, m.conf.Ip, m.conf.Port, m.conf.Database)
  126. m.db, err = sql.Open("taosSql", url)
  127. return err
  128. }
  129. func (m *taosSink) Collect(ctx api.StreamContext, item interface{}) error {
  130. logger := ctx.GetLogger()
  131. data, ok := item.([]byte)
  132. if !ok {
  133. logger.Debug("taos sink receive non string data")
  134. return nil
  135. }
  136. logger.Debugf("taos sink receive %s", item)
  137. var sliData []map[string]interface{}
  138. err := json.Unmarshal(data, &sliData)
  139. if nil != err {
  140. return err
  141. }
  142. for _, mapData := range sliData {
  143. sql, err := m.conf.buildSql(ctx, mapData)
  144. if nil != err {
  145. return err
  146. }
  147. logger.Debugf(sql)
  148. rows, err := m.db.Query(sql)
  149. if err != nil {
  150. return err
  151. }
  152. rows.Close()
  153. }
  154. return nil
  155. }
  156. func (m *taosSink) Close(ctx api.StreamContext) error {
  157. if m.db != nil {
  158. return m.db.Close()
  159. }
  160. return nil
  161. }
  162. func Taos() api.Sink {
  163. return &taosSink{}
  164. }