123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285 |
- package importer
- import (
- "encoding/base64"
- "encoding/json"
- "fmt"
- "strings"
- "sync"
- "sync/atomic"
- "time"
- "git.wecise.com/wecise/cgimport/odbc"
- "git.wecise.com/wecise/odb-go/odb"
- "git.wecise.com/wecise/util/cast"
- "git.wecise.com/wecise/util/cmap"
- "git.wecise.com/wecise/util/merrs"
- )
- type classdatainfo struct {
- *classinfo
- insertcount int64
- lastlogtime time.Time
- lastlogicount int64
- mutex sync.Mutex
- }
- var classdatainfos = cmap.New[string, *classdatainfo]()
- type ODBCImporter struct {
- client odb.Client
- }
- func NewODBCImporter() *ODBCImporter {
- odbci := &ODBCImporter{}
- if odbc.DevPhase&(odbc.DP_CREATECLASS|odbc.DP_INSERTDATA) != 0 {
- odbci.client = odbc.ODBC()
- }
- return odbci
- }
- // 根据数据修正类定义
- func (odbci *ODBCImporter) ReviseClassStruct() (err error) {
- for _, classname := range classnames {
- ci := classinfos.GetIFPresent(classname)
- if ci == nil {
- return merrs.NewError("classinfo not found " + classname)
- }
- _, err = classdatainfos.GetWithNew(classname, func() (cdi *classdatainfo, err error) {
- logger.Info("create class " + ci.classname)
- if odbci.client != nil {
- _, err = odbci.client.Query(ci.createmql).Do()
- if err != nil {
- return
- }
- }
- cdi = &classdatainfo{classinfo: ci}
- return
- })
- }
- if odbci.client != nil {
- for _, createedgemql := range createedgemqls {
- _, e := odbci.client.Query(createedgemql).Do()
- if e != nil && !strings.Contains(e.Error(), "already exist") {
- err = e
- return
- }
- }
- }
- return
- }
- func (odbci *ODBCImporter) InsertEdge(data map[string]any) (err error) {
- extraattr := map[string]string{}
- fromuid := ""
- touid := ""
- edgetype := ""
- for k, v := range data {
- switch k {
- case "FROMUNIQUEID":
- fromuid = cast.ToString(v)
- case "TOUNIQUEID":
- touid = cast.ToString(v)
- case "EDGETYPE":
- edgetype = cast.ToString(v)
- default:
- extraattr[k] = cast.ToString(v)
- }
- }
- if fromuid == "" {
- databs, _ := json.MarshalIndent(data, "", " ")
- return merrs.NewError("not found valid fromuniqueid in data ", merrs.SSMap{"data": string(databs)})
- }
- if touid == "" {
- databs, _ := json.MarshalIndent(data, "", " ")
- return merrs.NewError("not found valid touniqueid in data ", merrs.SSMap{"data": string(databs)})
- }
- if edgetype == "" {
- databs, _ := json.MarshalIndent(data, "", " ")
- return merrs.NewError("not found valid edgetype in data ", merrs.SSMap{"data": string(databs)})
- }
- edgetype = relations[edgetype]
- if edgetype == "" {
- databs, _ := json.MarshalIndent(data, "", " ")
- return merrs.NewError("not found valid edgetype in data ", merrs.SSMap{"data": string(databs)})
- }
- if odbci.client != nil {
- foid := get_object_id_from_cache(fromuid)
- toid := get_object_id_from_cache(touid)
- eabs, _ := json.Marshal(extraattr)
- quadmql := `quad "` + foid + `" ` + edgetype + ` + "` + toid + `" ` + string(eabs)
- _, err = odbci.client.Query(quadmql).Do()
- if err != nil {
- err = merrs.NewError(err, merrs.SSMaps{{"mql": quadmql}})
- logger.Error(err)
- return
- }
- }
- return
- }
- var cm_object_id_cache = cmap.New[string, chan string]()
- func object_id_cache(suid string) chan string {
- choid, _ := cm_object_id_cache.GetWithNew(suid,
- func() (chan string, error) {
- ch := make(chan string, 2)
- return ch, nil
- })
- return choid
- }
- func get_object_id_from_cache(suid string) string {
- choid := object_id_cache(suid)
- oid := <-choid
- push_object_id_into_cache(choid, oid)
- return oid
- }
- func push_object_id_into_cache(choid chan string, oid string) {
- choid <- oid
- if len(choid) == 2 {
- // 最多保留 1 个
- // chan cap = 2,第三个元素进不来
- // 进第二个元素的协程,清除第一个元素,允许其它协程后续进入新元素
- <-choid
- }
- }
- func object_id(classaliasname string, data map[string]any) (oid, suid string, err error) {
- uid := data["uniqueId"]
- if uid == nil {
- uid = data["UNIQUEID"]
- if uid == nil {
- databs, _ := json.MarshalIndent(data, "", " ")
- return "", "", merrs.NewError("not found uniqueid in data ", merrs.SSMap{"data": string(databs)})
- }
- }
- suid = cast.ToString(uid)
- if suid == "" {
- databs, _ := json.MarshalIndent(data, "", " ")
- return "", "", merrs.NewError("not found valid uniqueid in data ", merrs.SSMap{"data": string(databs)})
- }
- suid64 := base64.RawURLEncoding.EncodeToString([]byte(suid))
- return classaliasname + ":" + suid64, suid, nil
- }
- // 插入数据
- func (odbci *ODBCImporter) InsertData(classname string, data map[string]any) (err error) {
- cdi := classdatainfos.GetIFPresent(classname)
- if cdi == nil {
- return merrs.NewError("class not defined " + classname)
- }
- if cdi.insertmql == "" {
- return merrs.NewError("class no fields to insert " + classname)
- }
- oid, suid, e := object_id(cdi.aliasname, data)
- if e != nil {
- return e
- }
- values := []any{}
- for _, fn := range cdi.fieldslist {
- if fn == "id" {
- values = append(values, oid)
- continue
- }
- fi := cdi.fieldinfos[fn]
- if fi.datakey == "" {
- td := map[string]any{}
- for k, v := range data {
- if cdi.datakey_fieldinfos[k] == nil {
- td[k] = v
- }
- }
- tdbs, e := json.Marshal(td)
- if e != nil {
- return merrs.NewError(e)
- }
- values = append(values, string(tdbs))
- continue
- }
- v := data[fi.datakey]
- switch fi.fieldtype {
- case "set<varchar>":
- v = cast.ToStringSlice(v)
- case "timestamp":
- tv, e := cast.ToDateTimeE(v, "2006-01-02-15.04.05.000000")
- if e != nil {
- return merrs.NewError(fmt.Sprint("can't parse datetime value '", v, "'"))
- }
- v = tv.Format("2006-01-02 15:04:05.000000")
- }
- if fn == "tags" {
- v = append(cast.ToStringSlice(v), classname)
- }
- values = append(values, v)
- }
- if odbci.client != nil {
- _, err = odbci.client.Query(cdi.insertmql, values...).Do()
- if err != nil {
- err = merrs.NewError(err, merrs.SSMaps{{"mql": cdi.insertmql}, {"values": fmt.Sprint(values)}})
- logger.Error(err)
- return
- }
- push_object_id_into_cache(object_id_cache(suid), oid)
- }
- atomic.AddInt64(&cdi.insertcount, 1)
- cdi.mutex.Lock()
- if time.Since(cdi.lastlogtime) > 5*time.Second && cdi.lastlogicount != cdi.insertcount {
- cdi.lastlogtime = time.Now()
- cdi.lastlogicount = cdi.insertcount
- logger.Info("class", cdi.classname, "import", cdi.insertcount, "records")
- }
- cdi.mutex.Unlock()
- return
- }
- func (odbci *ODBCImporter) done() {
- classdatainfos.Fetch(func(cn string, cdi *classdatainfo) bool {
- cdi.mutex.Lock()
- if cdi.lastlogicount != cdi.insertcount {
- cdi.lastlogtime = time.Now()
- cdi.lastlogicount = cdi.insertcount
- logger.Info("class", cdi.classname, "import", cdi.insertcount, "records")
- }
- cdi.mutex.Unlock()
- return true
- })
- }
- func (odbci *ODBCImporter) alldone() {
- classdatainfos.Fetch(func(cn string, cdi *classdatainfo) bool {
- cdi.mutex.Lock()
- if cdi.insertcount != 0 {
- cdi.lastlogtime = time.Now()
- cdi.lastlogicount = cdi.insertcount
- logger.Info("class", cdi.classname, "import", cdi.insertcount, "records")
- }
- cdi.mutex.Unlock()
- return true
- })
- }
- func (odbci *ODBCImporter) reload() error {
- if odbci.client != nil {
- for i := len(classnames) - 1; i >= 0; i-- {
- classname := classnames[i]
- ci := classinfos.GetIFPresent(classname)
- if ci == nil {
- continue
- }
- for retry := 2; retry >= 0; retry-- {
- _, e := odbci.client.Query(`delete from "` + ci.classfullname + `" with version`).Do()
- _ = e
- _, e = odbci.client.Query(`drop class if exists "` + ci.classfullname + `"`).Do()
- if e != nil {
- if retry > 0 {
- continue
- }
- return e
- }
- }
- }
- }
- return nil
- }
|