tdengine.go 4.4 KB

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