模块整理,优化
This commit is contained in:
parent
b912c45092
commit
fee90b333d
@ -1,4 +1,4 @@
|
|||||||
package model
|
package dbservice
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"go_dreamfactory/modules"
|
"go_dreamfactory/modules"
|
||||||
@ -9,13 +9,13 @@ import (
|
|||||||
type Api_Comp struct {
|
type Api_Comp struct {
|
||||||
modules.MComp_GateComp
|
modules.MComp_GateComp
|
||||||
service core.IService
|
service core.IService
|
||||||
module *Model
|
module *DBService
|
||||||
}
|
}
|
||||||
|
|
||||||
func (this *Api_Comp) Init(service core.IService, module core.IModule, comp core.IModuleComp, options core.IModuleOptions) (err error) {
|
func (this *Api_Comp) Init(service core.IService, module core.IModule, comp core.IModuleComp, options core.IModuleOptions) (err error) {
|
||||||
this.MComp_GateComp.Init(service, module, comp, options)
|
this.MComp_GateComp.Init(service, module, comp, options)
|
||||||
this.service = service
|
this.service = service
|
||||||
this.module = module.(*Model)
|
this.module = module.(*DBService)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
@ -1,4 +1,4 @@
|
|||||||
package model
|
package dbservice
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"go_dreamfactory/lego/core"
|
"go_dreamfactory/lego/core"
|
@ -1,4 +1,4 @@
|
|||||||
package model
|
package dbservice
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
@ -14,6 +14,11 @@ import (
|
|||||||
|
|
||||||
const (
|
const (
|
||||||
WriteMaxNum uint32 = 1000 //一次性最处理条数
|
WriteMaxNum uint32 = 1000 //一次性最处理条数
|
||||||
|
ErrorMaxNum uint32 = 5 // 数据库操作最大失败次数
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
ErrorLogCount = make(map[string]uint32, 0)
|
||||||
)
|
)
|
||||||
|
|
||||||
type DB_Comp struct {
|
type DB_Comp struct {
|
||||||
@ -32,7 +37,7 @@ type QueryStruct struct {
|
|||||||
}
|
}
|
||||||
|
|
||||||
const (
|
const (
|
||||||
DB_ModelTable core.SqlTable = "model"
|
DB_ModelTable core.SqlTable = "model_log"
|
||||||
)
|
)
|
||||||
|
|
||||||
type IModel interface {
|
type IModel interface {
|
||||||
@ -53,13 +58,12 @@ func (this *DB_Comp) Model_UpdateDBByLog(uid string) (err error) {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
_delID := make([]string, 0) // 处理完成要删除的id
|
_delID := make([]string, 0) // 处理完成要删除的id
|
||||||
for _data.Next(context.TODO()) { // 处理删除逻辑
|
for _data.Next(context.TODO()) { // 处理删除逻辑
|
||||||
|
|
||||||
data := &comm.Autogenerated{}
|
data := &comm.Autogenerated{}
|
||||||
if err = _data.Decode(data); err == nil {
|
if err = _data.Decode(data); err != nil {
|
||||||
_delID = append(_delID, data.ID)
|
|
||||||
} else {
|
|
||||||
log.Errorf("Decode Data err : %v", err)
|
log.Errorf("Decode Data err : %v", err)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
@ -67,7 +71,7 @@ func (this *DB_Comp) Model_UpdateDBByLog(uid string) (err error) {
|
|||||||
if data.Act == string(comm.LogHandleType_Insert) {
|
if data.Act == string(comm.LogHandleType_Insert) {
|
||||||
if len(data.D) < 2 { // 参数校验
|
if len(data.D) < 2 { // 参数校验
|
||||||
log.Errorf("parameter len _id : %s,uid : %s d.len:%v", data.ID, data.UID, len(data.D))
|
log.Errorf("parameter len _id : %s,uid : %s d.len:%v", data.ID, data.UID, len(data.D))
|
||||||
break
|
continue
|
||||||
}
|
}
|
||||||
_obj := bson.M{}
|
_obj := bson.M{}
|
||||||
for _, v := range data.D[1].(bson.D) {
|
for _, v := range data.D[1].(bson.D) {
|
||||||
@ -78,29 +82,49 @@ func (this *DB_Comp) Model_UpdateDBByLog(uid string) (err error) {
|
|||||||
_, err := this.DB.InsertOne(core.SqlTable(_key), _obj)
|
_, err := this.DB.InsertOne(core.SqlTable(_key), _obj)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Errorf("insert %s db err:%v", (core.SqlTable(_key)), err)
|
log.Errorf("insert %s db err:%v", (core.SqlTable(_key)), err)
|
||||||
|
ErrorLogCount[data.ID]++
|
||||||
|
if ErrorLogCount[data.ID] >= ErrorMaxNum { // 实在是写失败了那就删除吧
|
||||||
|
log.Errorf("insert db err max num %s db err:%v", data.ID, err)
|
||||||
|
_, err = this.DB.DeleteOne(DB_ModelTable, bson.M{"_id": data.ID})
|
||||||
|
if err != nil {
|
||||||
|
log.Errorf("insert %s db err:%v", data.ID, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
continue
|
||||||
}
|
}
|
||||||
// _, err = this.DB.DeleteOne(DB_ModelTable, bson.M{"_id": data.ID})
|
// _, err = this.DB.DeleteOne(DB_ModelTable, bson.M{"_id": data.ID})
|
||||||
// if err != nil {
|
// if err != nil {
|
||||||
// log.Errorf("insert %s db err:%v", data.ID, err)
|
// log.Errorf("delete %s db err:%v", data.ID, err)
|
||||||
// }
|
// }
|
||||||
} else if data.Act == string(comm.LogHandleType_Delete) {
|
} else if data.Act == string(comm.LogHandleType_Delete) {
|
||||||
if len(data.D) < 2 { // 参数校验
|
if len(data.D) < 2 { // 参数校验
|
||||||
log.Errorf("parameter len _id : %s,uid : %s d.len:%v", data.ID, data.UID, len(data.D))
|
log.Errorf("parameter len _id : %s,uid : %s d.len:%v", data.ID, data.UID, len(data.D))
|
||||||
break
|
continue
|
||||||
}
|
}
|
||||||
|
_key := data.D[0].(string)
|
||||||
_objKey := make([]string, 0)
|
_objKey := make([]string, 0)
|
||||||
for _, v := range data.D[1].(bson.D) {
|
for _, v := range data.D[1].(bson.D) {
|
||||||
_objKey = append(_objKey, v.Value.(string))
|
_objKey = append(_objKey, v.Value.(string))
|
||||||
}
|
}
|
||||||
_, err = this.DB.DeleteMany(data.D[0].(core.SqlTable), bson.M{"_id": bson.M{"$in": _objKey}}, options.Delete())
|
_, err = this.DB.DeleteMany(core.SqlTable(_key), bson.M{"_id": bson.M{"$in": _objKey}}, options.Delete())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Errorf("delete %s db err:%v", data.D[0].(core.SqlTable), err)
|
log.Errorf("delete %s db err:%v", core.SqlTable(_key), err)
|
||||||
|
ErrorLogCount[data.ID]++
|
||||||
|
if ErrorLogCount[data.ID] >= ErrorMaxNum { // 实在是写失败了那就删除吧
|
||||||
|
log.Errorf("del db err max num %s db err:%v", data.ID, err)
|
||||||
|
_, err = this.DB.DeleteOne(DB_ModelTable, bson.M{"_id": data.ID})
|
||||||
|
if err != nil {
|
||||||
|
log.Errorf("insert %s db err:%v", data.ID, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
continue
|
||||||
}
|
}
|
||||||
} else { // update
|
} else { // update
|
||||||
if len(data.D) < 3 { // 参数校验
|
if len(data.D) < 3 { // 参数校验
|
||||||
log.Errorf("parameter len _id : %s,uid : %s d.len:%v", data.ID, data.UID, len(data.D))
|
log.Errorf("parameter len _id : %s,uid : %s d.len:%v", data.ID, data.UID, len(data.D))
|
||||||
break
|
continue
|
||||||
}
|
}
|
||||||
|
_key := data.D[0].(string)
|
||||||
where := data.D[1].(bson.D)
|
where := data.D[1].(bson.D)
|
||||||
_obj := &QueryStruct{}
|
_obj := &QueryStruct{}
|
||||||
for _, v := range where {
|
for _, v := range where {
|
||||||
@ -110,9 +134,22 @@ func (this *DB_Comp) Model_UpdateDBByLog(uid string) (err error) {
|
|||||||
for _, v := range query {
|
for _, v := range query {
|
||||||
_obj.Query[v.Key] = v
|
_obj.Query[v.Key] = v
|
||||||
}
|
}
|
||||||
this.DB.UpdateMany(data.D[0].(core.SqlTable), _obj.Selector, _obj.Query)
|
_, err := this.DB.UpdateMany(core.SqlTable(_key), _obj.Selector, _obj.Query)
|
||||||
|
if err != nil {
|
||||||
|
log.Errorf("Update %s db err:%v", core.SqlTable(_key), err)
|
||||||
|
ErrorLogCount[data.ID]++
|
||||||
|
if ErrorLogCount[data.ID] >= ErrorMaxNum { // 超过一定次数写失败了那就删除吧
|
||||||
|
log.Errorf("update db err max num %s db err:%v", data.ID, err)
|
||||||
|
_, err = this.DB.DeleteOne(DB_ModelTable, bson.M{"_id": data.ID})
|
||||||
|
if err != nil {
|
||||||
|
log.Errorf("insert %s db err:%v", data.ID, err)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
}
|
||||||
|
_delID = append(_delID, data.ID) // 都操作都成功了记录要删除的key
|
||||||
|
}
|
||||||
// 批量删除已处理的数据
|
// 批量删除已处理的数据
|
||||||
_, err = this.DB.DeleteMany(DB_ModelTable, bson.M{"_id": bson.M{"$in": _delID}}, options.Delete())
|
_, err = this.DB.DeleteMany(DB_ModelTable, bson.M{"_id": bson.M{"$in": _delID}}, options.Delete())
|
||||||
if err != nil {
|
if err != nil {
|
@ -1,4 +1,4 @@
|
|||||||
package model
|
package dbservice
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"go_dreamfactory/lego/core"
|
"go_dreamfactory/lego/core"
|
||||||
@ -9,12 +9,12 @@ import (
|
|||||||
type DBService_Comp struct {
|
type DBService_Comp struct {
|
||||||
cbase.ModuleCompBase
|
cbase.ModuleCompBase
|
||||||
task chan string
|
task chan string
|
||||||
module *Model
|
module *DBService
|
||||||
}
|
}
|
||||||
|
|
||||||
func (this *DBService_Comp) Init(service core.IService, module core.IModule, comp core.IModuleComp, options core.IModuleOptions) (err error) {
|
func (this *DBService_Comp) Init(service core.IService, module core.IModule, comp core.IModuleComp, options core.IModuleOptions) (err error) {
|
||||||
this.ModuleCompBase.Init(service, module, comp, options)
|
this.ModuleCompBase.Init(service, module, comp, options)
|
||||||
this.module = module.(*Model)
|
this.module = module.(*DBService)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
48
modules/dbservice/mail_test.go
Normal file
48
modules/dbservice/mail_test.go
Normal file
@ -0,0 +1,48 @@
|
|||||||
|
package dbservice
|
||||||
|
|
||||||
|
import (
|
||||||
|
"go_dreamfactory/comm"
|
||||||
|
"go_dreamfactory/lego/sys/log"
|
||||||
|
"go_dreamfactory/pb"
|
||||||
|
"os"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"go.mongodb.org/mongo-driver/bson/primitive"
|
||||||
|
)
|
||||||
|
|
||||||
|
var module = new(DBService)
|
||||||
|
|
||||||
|
func TestMain(m *testing.M) {
|
||||||
|
for i := 0; i < 50000; i++ {
|
||||||
|
//go func() {
|
||||||
|
_mail := &pb.DB_MailData{
|
||||||
|
ObjId: primitive.NewObjectID().Hex(),
|
||||||
|
UserId: "uid123",
|
||||||
|
Title: "系统邮件",
|
||||||
|
|
||||||
|
Contex: "恭喜获得专属礼包一份",
|
||||||
|
CreateTime: uint64(time.Now().Unix()),
|
||||||
|
DueTime: uint64(time.Now().Unix()) + 30*24*3600,
|
||||||
|
Check: false,
|
||||||
|
Reward: false,
|
||||||
|
}
|
||||||
|
//db.InsertModelLogs("mail", "uid123", _mail)
|
||||||
|
//InsertModelLogs("mail", "uid123", _mail)
|
||||||
|
data := &comm.Autogenerated{
|
||||||
|
ID: primitive.NewObjectID().Hex(),
|
||||||
|
UID: "uid123",
|
||||||
|
Act: string(comm.LogHandleType_Insert),
|
||||||
|
}
|
||||||
|
data.D = append(data.D, "mail") // D[0]
|
||||||
|
data.D = append(data.D, _mail) // D[1]
|
||||||
|
|
||||||
|
_, err1 := module.db_comp.DB.InsertOne("model_log", data)
|
||||||
|
if err1 != nil {
|
||||||
|
log.Errorf("insert model db err %v", err1)
|
||||||
|
}
|
||||||
|
//}()
|
||||||
|
}
|
||||||
|
time.Sleep(time.Second * 10)
|
||||||
|
defer os.Exit(m.Run())
|
||||||
|
}
|
@ -1,4 +1,4 @@
|
|||||||
package model
|
package dbservice
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"go_dreamfactory/comm"
|
"go_dreamfactory/comm"
|
||||||
@ -7,11 +7,11 @@ import (
|
|||||||
)
|
)
|
||||||
|
|
||||||
func NewModule() core.IModule {
|
func NewModule() core.IModule {
|
||||||
m := new(Model)
|
m := new(DBService)
|
||||||
return m
|
return m
|
||||||
}
|
}
|
||||||
|
|
||||||
type Model struct {
|
type DBService struct {
|
||||||
modules.ModuleBase
|
modules.ModuleBase
|
||||||
api_comp *Api_Comp
|
api_comp *Api_Comp
|
||||||
db_comp *DB_Comp
|
db_comp *DB_Comp
|
||||||
@ -19,17 +19,17 @@ type Model struct {
|
|||||||
configure_comp *Configure_Comp
|
configure_comp *Configure_Comp
|
||||||
}
|
}
|
||||||
|
|
||||||
func (this *Model) Init(service core.IService, module core.IModule, options core.IModuleOptions) (err error) {
|
func (this *DBService) Init(service core.IService, module core.IModule, options core.IModuleOptions) (err error) {
|
||||||
err = this.ModuleBase.Init(service, module, options)
|
err = this.ModuleBase.Init(service, module, options)
|
||||||
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
func (this *Model) GetType() core.M_Modules {
|
func (this *DBService) GetType() core.M_Modules {
|
||||||
return comm.SM_LogModelModule
|
return comm.SM_LogModelModule
|
||||||
}
|
}
|
||||||
|
|
||||||
func (this *Model) OnInstallComp() {
|
func (this *DBService) OnInstallComp() {
|
||||||
this.ModuleBase.OnInstallComp()
|
this.ModuleBase.OnInstallComp()
|
||||||
this.api_comp = this.RegisterComp(new(Api_Comp)).(*Api_Comp)
|
this.api_comp = this.RegisterComp(new(Api_Comp)).(*Api_Comp)
|
||||||
this.db_comp = this.RegisterComp(new(DB_Comp)).(*DB_Comp)
|
this.db_comp = this.RegisterComp(new(DB_Comp)).(*DB_Comp)
|
@ -1,25 +0,0 @@
|
|||||||
package model
|
|
||||||
|
|
||||||
import (
|
|
||||||
"go_dreamfactory/lego/sys/log"
|
|
||||||
"go_dreamfactory/pb"
|
|
||||||
"testing"
|
|
||||||
"time"
|
|
||||||
)
|
|
||||||
|
|
||||||
func TestCreatemoudles(t *testing.T) {
|
|
||||||
_mail := &pb.DB_MailData{
|
|
||||||
|
|
||||||
UserId: "uid123",
|
|
||||||
Title: "系统邮件",
|
|
||||||
|
|
||||||
Contex: "恭喜获得专属礼包一份",
|
|
||||||
CreateTime: uint64(time.Now().Unix()),
|
|
||||||
DueTime: uint64(time.Now().Unix()) + 30*24*3600,
|
|
||||||
Check: false,
|
|
||||||
Reward: false,
|
|
||||||
}
|
|
||||||
// obj.InsertModelLogs("mail", "uid123", _mail)
|
|
||||||
|
|
||||||
log.Debugf("insert : %v", _mail)
|
|
||||||
}
|
|
@ -13,10 +13,6 @@ import (
|
|||||||
"go.mongodb.org/mongo-driver/bson/primitive"
|
"go.mongodb.org/mongo-driver/bson/primitive"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
|
||||||
DB_ModelTable core.SqlTable = "model"
|
|
||||||
)
|
|
||||||
|
|
||||||
/*
|
/*
|
||||||
基础组件 缓存组件 读写缓存数据
|
基础组件 缓存组件 读写缓存数据
|
||||||
DB组件也封装进来
|
DB组件也封装进来
|
||||||
@ -27,6 +23,10 @@ type Model_Comp struct {
|
|||||||
DB mgo.ISys
|
DB mgo.ISys
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const (
|
||||||
|
DB_ModelTable core.SqlTable = "model_log"
|
||||||
|
)
|
||||||
|
|
||||||
//组件初始化接口
|
//组件初始化接口
|
||||||
func (this *Model_Comp) Init(service core.IService, module core.IModule, comp core.IModuleComp, options core.IModuleOptions) (err error) {
|
func (this *Model_Comp) Init(service core.IService, module core.IModule, comp core.IModuleComp, options core.IModuleOptions) (err error) {
|
||||||
this.ModuleCompBase.Init(service, module, comp, options)
|
this.ModuleCompBase.Init(service, module, comp, options)
|
||||||
|
@ -3,9 +3,9 @@ package main
|
|||||||
import (
|
import (
|
||||||
"flag"
|
"flag"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"go_dreamfactory/modules/dbservice"
|
||||||
"go_dreamfactory/modules/friend"
|
"go_dreamfactory/modules/friend"
|
||||||
"go_dreamfactory/modules/mail"
|
"go_dreamfactory/modules/mail"
|
||||||
"go_dreamfactory/modules/model"
|
|
||||||
"go_dreamfactory/modules/pack"
|
"go_dreamfactory/modules/pack"
|
||||||
"go_dreamfactory/modules/user"
|
"go_dreamfactory/modules/user"
|
||||||
"go_dreamfactory/services"
|
"go_dreamfactory/services"
|
||||||
@ -41,7 +41,7 @@ func main() {
|
|||||||
pack.NewModule(),
|
pack.NewModule(),
|
||||||
mail.NewModule(),
|
mail.NewModule(),
|
||||||
friend.NewModule(),
|
friend.NewModule(),
|
||||||
model.NewModule(),
|
dbservice.NewModule(),
|
||||||
)
|
)
|
||||||
|
|
||||||
}
|
}
|
||||||
|
2
sys/cache/init_test.go
vendored
2
sys/cache/init_test.go
vendored
@ -47,7 +47,7 @@ func TestMain(m *testing.M) {
|
|||||||
data.D = append(data.D, "mail") // D[0]
|
data.D = append(data.D, "mail") // D[0]
|
||||||
data.D = append(data.D, _mail) // D[1]
|
data.D = append(data.D, _mail) // D[1]
|
||||||
|
|
||||||
_, err1 := db.Defsys.Mgo().InsertOne("model", data)
|
_, err1 := db.Defsys.Mgo().InsertOne("model_log", data)
|
||||||
if err1 != nil {
|
if err1 != nil {
|
||||||
log.Errorf("insert model db err %v", err1)
|
log.Errorf("insert model db err %v", err1)
|
||||||
}
|
}
|
||||||
|
Loading…
Reference in New Issue
Block a user