odbcimporter.go 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242
  1. package importer
  2. import (
  3. "encoding/json"
  4. "strings"
  5. "sync"
  6. "sync/atomic"
  7. "time"
  8. "git.wecise.com/wecise/cgimport/odbc"
  9. "git.wecise.com/wecise/odb-go/odb"
  10. "git.wecise.com/wecise/util/cmap"
  11. "git.wecise.com/wecise/util/merrs"
  12. )
  13. type fieldinfo struct {
  14. fieldname string
  15. fieldtype string
  16. keyidx int // 主键顺序值,0为非主键
  17. datakey string // 对应数据中的键名
  18. }
  19. type classinfo struct {
  20. classname string
  21. nickname string
  22. fieldinfos map[string]*fieldinfo
  23. keyfields []string
  24. fieldslist []string
  25. insertmql string
  26. insertcount int64
  27. lastlogtime time.Time
  28. lastlogicount int64
  29. mutex sync.Mutex
  30. }
  31. type ODBCImporter struct {
  32. client odb.Client
  33. classinfos cmap.ConcurrentMap[string, *classinfo]
  34. }
  35. func NewODBCImporter() *ODBCImporter {
  36. odbci := &ODBCImporter{
  37. classinfos: cmap.New[string, *classinfo](),
  38. }
  39. if odbc.DevPhase&(odbc.DP_CREATECLASS|odbc.DP_INSERTDATA) != 0 {
  40. odbci.client = odbc.ODBC()
  41. }
  42. return odbci
  43. }
  44. // 根据数据修正类定义
  45. func (odbci *ODBCImporter) ReviseClassStruct(data map[string]any) (classname string, err error) {
  46. switch {
  47. case data["uniqueId"] != nil:
  48. classname = "/cgitest/x10/x1002"
  49. case data["UNIQUEID"] != nil:
  50. classname = "/cgitest/x10/x1001"
  51. case data["FROMUNIQUEID"] != nil:
  52. classname = "/cgitest/x10/x1003"
  53. default:
  54. bs, e := json.MarshalIndent(data, "", " ")
  55. if e != nil {
  56. err = e
  57. return
  58. }
  59. err = merrs.NewError("no mapping classname", merrs.SSMaps{{"data": string(bs)}})
  60. }
  61. err = odbci.createClass("cgitest", "/cgitest", nil)
  62. if err != nil {
  63. return
  64. }
  65. err = odbci.createClass("x10", "/cgitest/x10", nil)
  66. if err != nil {
  67. return
  68. }
  69. err = odbci.createClass("x1001", "/cgitest/x10/x1001", []*fieldinfo{
  70. {fieldname: "uniqueid", datakey: "UNIQUEID", fieldtype: "varchar", keyidx: 1},
  71. {fieldname: "distname", datakey: "BASENAME", fieldtype: "varchar"},
  72. })
  73. if err != nil {
  74. return
  75. }
  76. err = odbci.createClass("x1002", "/cgitest/x10/x1002", []*fieldinfo{
  77. {fieldname: "uniqueid", datakey: "uniqueId", fieldtype: "varchar", keyidx: 1},
  78. {fieldname: "distname", datakey: "distName", fieldtype: "varchar"},
  79. })
  80. if err != nil {
  81. return
  82. }
  83. err = odbci.createClass("x1003", "/cgitest/x10/x1003", []*fieldinfo{
  84. {fieldname: "fromuniqueid", datakey: "FROMUNIQUEID", fieldtype: "varchar", keyidx: 1},
  85. {fieldname: "touniqueid", datakey: "TOUNIQUEID", fieldtype: "varchar", keyidx: 2},
  86. })
  87. if err != nil {
  88. return
  89. }
  90. return
  91. }
  92. // 插入数据
  93. func (odbci *ODBCImporter) InsertData(classname string, data map[string]any) (err error) {
  94. ci := odbci.classinfos.GetIFPresent(classname)
  95. if ci == nil {
  96. return merrs.NewError("class not defined " + classname)
  97. }
  98. if ci.insertmql == "" {
  99. return merrs.NewError("class no fields to insert " + classname)
  100. }
  101. values := []any{}
  102. for _, fn := range ci.fieldslist {
  103. fi := ci.fieldinfos[fn]
  104. v := data[fi.datakey]
  105. values = append(values, v)
  106. }
  107. if odbci.client != nil {
  108. _, err = odbci.client.Query(ci.insertmql, values...).Do()
  109. if err != nil {
  110. return
  111. }
  112. }
  113. atomic.AddInt64(&ci.insertcount, 1)
  114. ci.mutex.Lock()
  115. if time.Since(ci.lastlogtime) > 5*time.Second && ci.lastlogicount != ci.insertcount {
  116. ci.lastlogtime = time.Now()
  117. ci.lastlogicount = ci.insertcount
  118. logger.Info("class", ci.classname, "import", ci.insertcount, "records")
  119. }
  120. ci.mutex.Unlock()
  121. return
  122. }
  123. func (odbci *ODBCImporter) done() {
  124. odbci.classinfos.Fetch(func(cn string, ci *classinfo) bool {
  125. ci.mutex.Lock()
  126. if ci.lastlogicount != ci.insertcount {
  127. ci.lastlogtime = time.Now()
  128. ci.lastlogicount = ci.insertcount
  129. logger.Info("class", ci.classname, "import", ci.insertcount, "records")
  130. }
  131. ci.mutex.Unlock()
  132. return true
  133. })
  134. }
  135. func (odbci *ODBCImporter) alldone() {
  136. odbci.classinfos.Fetch(func(cn string, ci *classinfo) bool {
  137. ci.mutex.Lock()
  138. if ci.insertcount != 0 {
  139. ci.lastlogtime = time.Now()
  140. ci.lastlogicount = ci.insertcount
  141. logger.Info("class", ci.classname, "import", ci.insertcount, "records")
  142. }
  143. ci.mutex.Unlock()
  144. return true
  145. })
  146. }
  147. func (odbci *ODBCImporter) reload() error {
  148. if odbci.client != nil {
  149. _, e := odbci.client.Query(`delete from "/cgitest/x10/x1001" with version`).Do()
  150. _ = e
  151. _, e = odbci.client.Query(`delete from "/cgitest/x10/x1002" with version`).Do()
  152. _ = e
  153. _, e = odbci.client.Query(`delete from "/cgitest/x10/x1003" with version`).Do()
  154. _ = e
  155. _, e = odbci.client.Query(`drop class if exists "/cgitest/x10/x1001"`).Do()
  156. if e != nil {
  157. return e
  158. }
  159. _, e = odbci.client.Query(`drop class if exists "/cgitest/x10/x1002"`).Do()
  160. if e != nil {
  161. return e
  162. }
  163. _, e = odbci.client.Query(`drop class if exists "/cgitest/x10/x1003"`).Do()
  164. if e != nil {
  165. return e
  166. }
  167. _, e = odbci.client.Query(`drop class if exists "/cgitest/x10"`).Do()
  168. if e != nil {
  169. return e
  170. }
  171. }
  172. return nil
  173. }
  174. // 新建类
  175. func (odbci *ODBCImporter) createClass(nickname, classname string, fieldinfoslist []*fieldinfo) (err error) {
  176. _, err = odbci.classinfos.GetWithNew(classname, func() (ci *classinfo, err error) {
  177. logger.Info("create class " + classname)
  178. fieldinfos := map[string]*fieldinfo{}
  179. fieldslist := []string{}
  180. keyfields := []string{}
  181. mql := `create class if not exists ` + classname + `(`
  182. if len(fieldinfoslist) > 0 {
  183. field_defines := []string{}
  184. for _, fi := range fieldinfoslist {
  185. field_defines = append(field_defines, fi.fieldname+" "+fi.fieldtype)
  186. fieldslist = append(fieldslist, fi.fieldname)
  187. if fi.keyidx > 0 {
  188. for len(keyfields) < fi.keyidx {
  189. keyfields = append(keyfields, "")
  190. }
  191. keyfields[fi.keyidx-1] = fi.fieldname
  192. }
  193. fieldinfos[fi.fieldname] = fi
  194. }
  195. mql += strings.Join(field_defines, ",")
  196. }
  197. if len(keyfields) > 0 {
  198. mql += ", keys(" + strings.Join(keyfields, ",") + ")"
  199. }
  200. mql += `)with namespace="cgitest" and key=manu and nickname='` + nickname + `'`
  201. if odbci.client != nil {
  202. _, err = odbci.client.Query(mql).Do()
  203. if err != nil {
  204. return
  205. }
  206. }
  207. var insertmql string
  208. if len(fieldslist) > 0 {
  209. insertmql = `insert into ` + classname + "(" + strings.Join(fieldslist, ",") + ")values(" + strings.Repeat(",?", len(fieldslist))[1:] + ")"
  210. }
  211. ci = &classinfo{
  212. classname: classname,
  213. nickname: nickname,
  214. fieldinfos: fieldinfos,
  215. keyfields: keyfields,
  216. fieldslist: fieldslist,
  217. insertmql: insertmql,
  218. }
  219. return
  220. })
  221. return
  222. }
  223. // 修改类定义
  224. func (odbci *ODBCImporter) alterClass() (err error) {
  225. return
  226. }