diff --git a/cmd/climc/shell/events/splitable.go b/cmd/climc/shell/events/splitable.go new file mode 100644 index 0000000000..e76e26ab6a --- /dev/null +++ b/cmd/climc/shell/events/splitable.go @@ -0,0 +1,70 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package events + +import ( + "yunion.io/x/jsonutils" + + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/modulebase" + "yunion.io/x/onecloud/pkg/mcclient/modules" +) + +func init() { + type EventSplitableOptions struct { + Service string `help:"service" choices:"compute|identity|image" default:"compute"` + } + R(&EventSplitableOptions{}, "logs-splitable", "Show splitable info of event table", func(s *mcclient.ClientSession, args *EventSplitableOptions) error { + var results jsonutils.JSONObject + var err error + switch args.Service { + case "identity": + results, err = modules.IdentityLogs.Get(s, "splitable", nil) + case "image": + results, err = modules.ImageLogs.Get(s, "splitable", nil) + default: + results, err = modules.Logs.Get(s, "splitable", nil) + } + if err != nil { + return err + } + tables, err := results.GetArray() + if err != nil { + return err + } + listResult := &modulebase.ListResult{ + Data: tables, + } + printList(listResult, nil) + return nil + }) + R(&EventSplitableOptions{}, "logs-purge", "Purge obsolete splitable of event table", func(s *mcclient.ClientSession, args *EventSplitableOptions) error { + var results jsonutils.JSONObject + var err error + switch args.Service { + case "identity": + results, err = modules.IdentityLogs.PerformClassAction(s, "purge-splitable", nil) + case "image": + results, err = modules.ImageLogs.PerformClassAction(s, "purge-splitable", nil) + default: + results, err = modules.Logs.PerformClassAction(s, "purge-splitable", nil) + } + if err != nil { + return err + } + printObject(results) + return nil + }) +} diff --git a/pkg/apis/identity/consts.go b/pkg/apis/identity/consts.go index 02503e7f1b..3de0a7d23a 100644 --- a/pkg/apis/identity/consts.go +++ b/pkg/apis/identity/consts.go @@ -163,6 +163,8 @@ var ( "etcd_cacert", "etcd_cert", "etcd_key", + "splitable_max_duration_hours", + "splitable_max_keep_segments", // ############################ // keystone blacklist options diff --git a/pkg/cloudcommon/consts/opslog.go b/pkg/cloudcommon/consts/opslog.go index c4b9930666..05f9450f1c 100644 --- a/pkg/cloudcommon/consts/opslog.go +++ b/pkg/cloudcommon/consts/opslog.go @@ -14,8 +14,12 @@ package consts +import "time" + var ( - globalOpsLogEnabled = true + globalOpsLogEnabled = true + splitableMaxDurationHours = 24 * 30 // 30 days + splitableMaxKeepSegments = 6 // 6 * 30 days, half year ) func DisableOpsLog() { @@ -25,3 +29,19 @@ func DisableOpsLog() { func OpsLogEnabled() bool { return globalOpsLogEnabled } + +func SetSplitableMaxKeepSegments(cnt int) { + splitableMaxKeepSegments = cnt +} + +func SetSplitableMaxDurationHours(h int) { + splitableMaxDurationHours = h +} + +func SplitableMaxKeepSegments() int { + return splitableMaxKeepSegments +} + +func SplitableMaxDuration() time.Duration { + return time.Hour * time.Duration(splitableMaxDurationHours) +} diff --git a/pkg/cloudcommon/db/db_dispatcher.go b/pkg/cloudcommon/db/db_dispatcher.go index 29e646c842..5f766bc696 100644 --- a/pkg/cloudcommon/db/db_dispatcher.go +++ b/pkg/cloudcommon/db/db_dispatcher.go @@ -536,11 +536,6 @@ func ListItems(manager IModelManager, ctx context.Context, userCred mcclient.Tok useRawQuery = true } } - if useRawQuery { - q = manager.RawQuery() - } else { - q = manager.Query() - } queryDict, ok := query.(*jsonutils.JSONDict) if !ok { @@ -559,17 +554,81 @@ func ListItems(manager IModelManager, ctx context.Context, userCred mcclient.Tok return nil, err } + pagingConf := manager.GetPagingConfig() + if pagingConf == nil { + if limit <= 0 { + limit = consts.GetDefaultPagingLimit() + } + } else { + if limit <= 0 { + limit = int64(pagingConf.DefaultLimit) + } + } + + splitable := manager.GetSplitTable() + if splitable != nil { + // handle splitable query, query each subtable, then union results + metas, err := splitable.GetTableMetas() + if err != nil { + return nil, errors.Wrap(err, "splitable.GetTableMetas") + } + var subqs []sqlchemy.IQuery + for _, meta := range metas { + ts := splitable.GetTableSpec(meta) + subq := ts.Query() + subq, err = listItemQueryFiltersRaw(manager, ctx, subq, userCred, queryDict, policy.PolicyActionList, true, useRawQuery) + if err != nil { + return nil, errors.Wrap(err, "listItemQueryFiltersRaw") + } + if pagingConf != nil { + if limit > 0 { + subq = subq.Limit(int(limit) + 1) + } + if len(pagingMarker) > 0 { + markers := decodePagingMarker(pagingMarker) + for markerIdx, marker := range markers { + if markerIdx < len(pagingConf.MarkerFields) { + if pagingConf.Order == sqlchemy.SQL_ORDER_ASC { + subq = subq.GE(pagingConf.MarkerFields[markerIdx], marker) + } else { + subq = subq.LE(pagingConf.MarkerFields[markerIdx], marker) + } + } + } + } + for _, f := range pagingConf.MarkerFields { + if pagingConf.Order == sqlchemy.SQL_ORDER_ASC { + subq = subq.Asc(f) + } else { + subq = subq.Desc(f) + } + } + } + if limit > 0 { + subq = subq.Limit(int(limit) + 1) + } + subqs = append(subqs, subq) + } + union, err := sqlchemy.UnionWithError(subqs...) + if err != nil { + return nil, errors.Wrap(err, "sqlchemy.UnionWithError") + } + q = union.Query() + } else { + if useRawQuery { + q = manager.RawQuery() + } else { + q = manager.Query() + } + } + q, err = listItemQueryFiltersRaw(manager, ctx, q, userCred, queryDict, policy.PolicyActionList, true, useRawQuery) if err != nil { - return nil, err + return nil, errors.Wrap(err, "listItemQueryFiltersRaw") } var totalCnt int - pagingConf := manager.GetPagingConfig() if pagingConf == nil { - if limit == 0 { - limit = consts.GetDefaultPagingLimit() - } totalCnt, err = q.CountWithError() if err != nil { return nil, err @@ -579,10 +638,6 @@ func ListItems(manager IModelManager, ctx context.Context, userCred mcclient.Tok emptyList := modulebase.ListResult{Data: []jsonutils.JSONObject{}} return &emptyList, nil } - } else { - if limit <= 0 { - limit = int64(pagingConf.DefaultLimit) - } } if int64(totalCnt) > maxLimit && (limit <= 0 || limit > maxLimit) { limit = maxLimit diff --git a/pkg/cloudcommon/db/interface.go b/pkg/cloudcommon/db/interface.go index f0aa8b1634..b959b5bbd9 100644 --- a/pkg/cloudcommon/db/interface.go +++ b/pkg/cloudcommon/db/interface.go @@ -27,6 +27,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/object" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/util/rbacutils" + "yunion.io/x/onecloud/pkg/util/splitable" "yunion.io/x/onecloud/pkg/util/stringutils2" ) @@ -130,6 +131,8 @@ type IModelManager interface { GetPagingConfig() *SPagingConfig GetI18N(ctx context.Context, idstr string, resObj jsonutils.JSONObject) *jsonutils.JSONDict + + GetSplitTable() *splitable.SSplitTableSpec } type IModel interface { diff --git a/pkg/cloudcommon/db/modelbase.go b/pkg/cloudcommon/db/modelbase.go index 95ed2047e1..bf0d7d1879 100644 --- a/pkg/cloudcommon/db/modelbase.go +++ b/pkg/cloudcommon/db/modelbase.go @@ -32,6 +32,7 @@ import ( "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/util/rbacutils" + "yunion.io/x/onecloud/pkg/util/splitable" "yunion.io/x/onecloud/pkg/util/stringutils2" ) @@ -52,7 +53,11 @@ type SModelBaseManager struct { } func NewModelBaseManager(model interface{}, tableName string, keyword string, keywordPlural string) SModelBaseManager { - ts := newTableSpec(model, tableName) + return NewModelBaseManagerWithSplitable(model, tableName, keyword, keywordPlural, "", "", 0, 0) +} + +func NewModelBaseManagerWithSplitable(model interface{}, tableName string, keyword string, keywordPlural string, indexField string, dateField string, maxDuration time.Duration, maxSegments int) SModelBaseManager { + ts := newTableSpec(model, tableName, indexField, dateField, maxDuration, maxSegments) modelMan := SModelBaseManager{tableSpec: ts, keyword: keyword, keywordPlural: keywordPlural} return modelMan } @@ -82,6 +87,10 @@ func (manager *SModelBaseManager) TableSpec() ITableSpec { return manager.tableSpec } +func (manager *SModelBaseManager) GetSplitTable() *splitable.SSplitTableSpec { + return manager.tableSpec.GetSplitTable() +} + func (manager *SModelBaseManager) Keyword() string { return manager.keyword } @@ -433,6 +442,38 @@ func (manager *SModelBaseManager) GetI18N(ctx context.Context, idstr string, res return nil } +func (manager *SModelBaseManager) AllowGetPropertySplitable(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { + return true +} + +func (manager *SModelBaseManager) GetPropertySplitable(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (jsonutils.JSONObject, error) { + splitable := manager.GetSplitTable() + if splitable == nil { + return nil, errors.Wrap(httperrors.ErrNotSupported, "not splitable") + } + metas, err := splitable.GetTableMetas() + if err != nil { + return nil, errors.Wrap(err, "GetTableMetas") + } + return jsonutils.Marshal(metas), nil +} + +func (manager *SModelBaseManager) AllowPerformPurgeSplitable(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return true +} + +func (manager *SModelBaseManager) PerformPurgeSplitable(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + splitable := manager.GetSplitTable() + if splitable == nil { + return nil, errors.Wrap(httperrors.ErrNotSupported, "not splitable") + } + err := splitable.Purge() + if err != nil { + return nil, errors.Wrap(err, "Purge") + } + return nil, nil +} + func (model *SModelBase) GetId() string { return "" } diff --git a/pkg/cloudcommon/db/opslog.go b/pkg/cloudcommon/db/opslog.go index 2870d990c2..7b76135d49 100644 --- a/pkg/cloudcommon/db/opslog.go +++ b/pkg/cloudcommon/db/opslog.go @@ -78,7 +78,16 @@ var _ IModel = (*SOpsLog)(nil) var opslogQueryWorkerMan *appsrv.SWorkerManager func init() { - OpsLog = &SOpsLogManager{NewModelBaseManager(SOpsLog{}, "opslog_tbl", "event", "events")} + OpsLog = &SOpsLogManager{NewModelBaseManagerWithSplitable( + SOpsLog{}, + "opslog_tbl", + "event", + "events", + "id", + "ops_time", + consts.SplitableMaxDuration(), + consts.SplitableMaxKeepSegments(), + )} OpsLog.SetVirtualObject(OpsLog) opslogQueryWorkerMan = appsrv.NewWorkerManager("opslog_query_worker", 2, 1024, true) diff --git a/pkg/cloudcommon/db/tablespec.go b/pkg/cloudcommon/db/tablespec.go index 9355922818..264c3f2f34 100644 --- a/pkg/cloudcommon/db/tablespec.go +++ b/pkg/cloudcommon/db/tablespec.go @@ -17,6 +17,7 @@ package db import ( "context" "reflect" + "time" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -25,6 +26,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/informer" "yunion.io/x/onecloud/pkg/util/nopanic" + "yunion.io/x/onecloud/pkg/util/splitable" ) type ITableSpec interface { @@ -32,29 +34,52 @@ type ITableSpec interface { Columns() []sqlchemy.IColumnSpec PrimaryColumns() []sqlchemy.IColumnSpec DataType() reflect.Type - CreateSQL() string + // CreateSQL() string Instance() *sqlchemy.STable ColumnSpec(name string) sqlchemy.IColumnSpec Insert(ctx context.Context, dt interface{}) error InsertOrUpdate(ctx context.Context, dt interface{}) error Update(ctx context.Context, dt interface{}, doUpdate func() error) (sqlchemy.UpdateDiffs, error) Fetch(dt interface{}) error - FetchAll(dest interface{}) error + // FetchAll(dest interface{}) error SyncSQL() []string DropForeignKeySQL() []string AddIndex(unique bool, cols ...string) bool Increment(ctx context.Context, diff interface{}, target interface{}) error Decrement(ctx context.Context, diff interface{}, target interface{}) error + + GetSplitTable() *splitable.SSplitTableSpec } type sTableSpec struct { - *sqlchemy.STableSpec + sqlchemy.ITableSpec } -func newTableSpec(model interface{}, tableName string) ITableSpec { - return &sTableSpec{ - STableSpec: sqlchemy.NewTableSpecFromStruct(model, tableName), +func newTableSpec(model interface{}, tableName string, indexField string, dateField string, maxDuration time.Duration, maxSegments int) ITableSpec { + var itbl sqlchemy.ITableSpec + if len(indexField) > 0 && len(dateField) > 0 { + var err error + itbl, err = splitable.NewSplitTableSpec(model, tableName, indexField, dateField, maxDuration, maxSegments) + if err != nil { + log.Errorf("NewSplitTableSpec %s %s", tableName, err) + return nil + } else { + log.Debugf("table %s maxDuration %d hour maxSegements %d", tableName, maxDuration/time.Hour, maxSegments) + } + } else { + itbl = sqlchemy.NewTableSpecFromStruct(model, tableName) } + return &sTableSpec{ + ITableSpec: itbl, + } +} + +func (ts *sTableSpec) GetSplitTable() *splitable.SSplitTableSpec { + sts, ok := ts.ITableSpec.(*splitable.SSplitTableSpec) + if ok { + return sts + } + return nil } func (ts *sTableSpec) newInformerModel(dt interface{}) (*informer.ModelObject, error) { @@ -91,7 +116,7 @@ func (ts *sTableSpec) isMarkDeleted(dt interface{}) (bool, error) { } func (ts *sTableSpec) Insert(ctx context.Context, dt interface{}) error { - if err := ts.STableSpec.Insert(dt); err != nil { + if err := ts.ITableSpec.Insert(dt); err != nil { return err } ts.inform(ctx, dt, informer.Create) @@ -99,7 +124,7 @@ func (ts *sTableSpec) Insert(ctx context.Context, dt interface{}) error { } func (ts *sTableSpec) InsertOrUpdate(ctx context.Context, dt interface{}) error { - if err := ts.STableSpec.InsertOrUpdate(dt); err != nil { + if err := ts.ITableSpec.InsertOrUpdate(dt); err != nil { return err } ts.inform(ctx, dt, informer.Create) @@ -108,7 +133,7 @@ func (ts *sTableSpec) InsertOrUpdate(ctx context.Context, dt interface{}) error func (ts *sTableSpec) Update(ctx context.Context, dt interface{}, doUpdate func() error) (sqlchemy.UpdateDiffs, error) { oldObj := jsonutils.Marshal(dt) - diffs, err := ts.STableSpec.Update(dt, doUpdate) + diffs, err := ts.ITableSpec.Update(dt, doUpdate) if err != nil { return nil, err } @@ -130,7 +155,7 @@ func (ts *sTableSpec) Update(ctx context.Context, dt interface{}, doUpdate func( func (ts *sTableSpec) Increment(ctx context.Context, diff, target interface{}) error { oldObj := jsonutils.Marshal(target) - err := ts.STableSpec.Increment(diff, target) + err := ts.ITableSpec.Increment(diff, target) if err != nil { return errors.Wrap(err, "Increment") } @@ -140,7 +165,7 @@ func (ts *sTableSpec) Increment(ctx context.Context, diff, target interface{}) e func (ts *sTableSpec) Decrement(ctx context.Context, diff, target interface{}) error { oldObj := jsonutils.Marshal(target) - err := ts.STableSpec.Decrement(diff, target) + err := ts.ITableSpec.Decrement(diff, target) if err != nil { return err } diff --git a/pkg/cloudcommon/options/options.go b/pkg/cloudcommon/options/options.go index 1d0ca0d90a..03de65563d 100644 --- a/pkg/cloudcommon/options/options.go +++ b/pkg/cloudcommon/options/options.go @@ -138,6 +138,9 @@ type DBOptions struct { LockmanMethod string `help:"method for lock synchronization" choices:"inmemory|etcd" default:"inmemory"` + // SplitableMaxKeepSegments int `help:"maximal segements of splitable to keep, default 6 segments" default:"6"` + // SplitableMaxDurationHours int `help:"maximal number of hours that a splitable segement lasts, default 30 days" default:"720"` + EtcdOptions EtcdLockPrefix string `help:"prefix of etcd lock records" default:"/onecloud/lockman"` diff --git a/pkg/cloudevent/models/cloudevents.go b/pkg/cloudevent/models/cloudevents.go index 298c604d36..6d0f9786e5 100644 --- a/pkg/cloudevent/models/cloudevents.go +++ b/pkg/cloudevent/models/cloudevents.go @@ -24,6 +24,7 @@ import ( "yunion.io/x/sqlchemy" api "yunion.io/x/onecloud/pkg/apis/cloudevent" + "yunion.io/x/onecloud/pkg/cloudcommon/consts" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/mcclient" @@ -40,11 +41,15 @@ var CloudeventManager *SCloudeventManager func init() { CloudeventManager = &SCloudeventManager{ - SModelBaseManager: db.NewModelBaseManager( + SModelBaseManager: db.NewModelBaseManagerWithSplitable( SCloudevent{}, "cloudevents_tbl", "cloudevent", "cloudevents", + "event_id", + "created_at", + consts.SplitableMaxDuration(), + consts.SplitableMaxKeepSegments(), ), } CloudeventManager.SetVirtualObject(CloudeventManager) diff --git a/pkg/logger/models/actionlog.go b/pkg/logger/models/actionlog.go index fa60126ba8..5e8c14cb19 100644 --- a/pkg/logger/models/actionlog.go +++ b/pkg/logger/models/actionlog.go @@ -21,6 +21,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/cloudcommon/consts" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/auth" @@ -59,7 +60,16 @@ func init() { InitActionWhiteList() ActionLog = &SActionlogManager{ SOpsLogManager: db.SOpsLogManager{ - SModelBaseManager: db.NewModelBaseManager(SActionlog{}, "action_tbl", "action", "actions"), + SModelBaseManager: db.NewModelBaseManagerWithSplitable( + SActionlog{}, + "action_tbl", + "action", + "actions", + "id", + "start_time", + consts.SplitableMaxDuration(), + consts.SplitableMaxKeepSegments(), + ), }, } ActionLog.SetVirtualObject(ActionLog) diff --git a/pkg/logger/models/baremetalevents.go b/pkg/logger/models/baremetalevents.go index d26b28aa80..d6d18320b0 100644 --- a/pkg/logger/models/baremetalevents.go +++ b/pkg/logger/models/baremetalevents.go @@ -23,6 +23,7 @@ import ( "yunion.io/x/sqlchemy" api "yunion.io/x/onecloud/pkg/apis/logger" + "yunion.io/x/onecloud/pkg/cloudcommon/consts" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/util/rbacutils" @@ -51,11 +52,15 @@ var BaremetalEventManager *SBaremetalEventManager func init() { BaremetalEventManager = &SBaremetalEventManager{ - SModelBaseManager: db.NewModelBaseManager( + SModelBaseManager: db.NewModelBaseManagerWithSplitable( SBaremetalEvent{}, "baremetal_event_tbl", "baremetalevent", "baremetalevents", + "id", + "created", + consts.SplitableMaxDuration(), + consts.SplitableMaxKeepSegments(), ), } BaremetalEventManager.SetVirtualObject(BaremetalEventManager) diff --git a/pkg/util/splitable/doc.go b/pkg/util/splitable/doc.go new file mode 100644 index 0000000000..963cb2cd9a --- /dev/null +++ b/pkg/util/splitable/doc.go @@ -0,0 +1,15 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package splitable // import "yunion.io/x/onecloud/pkg/util/splitable" diff --git a/pkg/util/splitable/insert.go b/pkg/util/splitable/insert.go new file mode 100644 index 0000000000..58bd0196fc --- /dev/null +++ b/pkg/util/splitable/insert.go @@ -0,0 +1,92 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package splitable + +import ( + "fmt" + "reflect" + "time" + + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/util/reflectutils" + "yunion.io/x/sqlchemy" +) + +func (t *SSplitTableSpec) Insert(dt interface{}) error { + metas, err := t.GetTableMetas() + if err != nil { + return errors.Wrap(err, "GetTableMeta") + } + var lastDate time.Time + vs := reflectutils.FetchAllStructFieldValueSet(reflect.Indirect(reflect.ValueOf(dt))) + if lastDateV, ok := vs.GetValue(t.dateField); !ok { + return errors.Wrap(errors.ErrInvalidStatus, "no dateField found") + } else { + lastDate = lastDateV.Interface().(time.Time) + } + + var lastRecIndex int64 + var lastRecDate time.Time + var lastTableSpec *sqlchemy.STableSpec + + newMeta := false + if len(metas) > 0 { + lastMeta := metas[len(metas)-1] + if lastDate.Sub(lastMeta.StartDate) > t.maxDuration { + lastTable := t.GetTableSpec(lastMeta) + ti := lastTable.Instance() + q := ti.Query(sqlchemy.MAX("last_index", ti.Field(t.indexField)), sqlchemy.MAX("last_date", ti.Field(t.dateField))) + r := q.Row() + err := r.Scan(&lastRecIndex, &lastRecDate) + if err != nil { + return errors.Wrap(err, "scan lastRecIndex and lastRecDate") + } + // seal last meta + _, err = t.metaSpec.Update(&lastMeta, func() error { + lastMeta.End = lastRecIndex + lastMeta.EndDate = lastRecDate + return nil + }) + if err != nil { + return errors.Wrap(err, "Update last meta") + } + newMeta = true + } else { + lastTableSpec = t.GetTableSpec(lastMeta) + } + } else { + newMeta = true + } + if newMeta { + // insert a new metadata + meta := STableMetadata{ + Table: fmt.Sprintf("%s_%d", t.tableName, lastDate.Unix()), + Start: lastRecIndex + 1, + StartDate: lastDate, + } + err := t.metaSpec.Insert(&meta) + if err != nil { + return errors.Wrap(err, "insert new meta") + } + // create new table + newTable := t.GetTableSpec(meta) + err = newTable.Sync() + if err != nil { + return errors.Wrap(err, "sync new table") + } + lastTableSpec = newTable + } + return lastTableSpec.Insert(dt) +} diff --git a/pkg/util/splitable/metadata.go b/pkg/util/splitable/metadata.go new file mode 100644 index 0000000000..b7d38a0425 --- /dev/null +++ b/pkg/util/splitable/metadata.go @@ -0,0 +1,50 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package splitable + +import ( + "database/sql" + "time" + + "yunion.io/x/pkg/errors" + "yunion.io/x/sqlchemy" +) + +type STableMetadata struct { + Id int64 `primary:"true" auto_increment:"true"` + Table string `width:"64" charset:"ascii"` + Start int64 `nullable:"true"` + End int64 `nullable:"true"` + StartDate time.Time `nullable:"true"` + EndDate time.Time `nullable:"true"` + Deleted bool `nullable:"false"` + DeleteAt time.Time `nullable:"true"` + CreatedAt time.Time `nullable:"false" created_at:"true"` +} + +func (spec *SSplitTableSpec) GetTableMetas() ([]STableMetadata, error) { + q := spec.metaSpec.Query().Asc("id").IsFalse("deleted") + metas := make([]STableMetadata, 0) + err := q.All(&metas) + if err != nil && errors.Cause(err) != sql.ErrNoRows { + return nil, errors.Wrap(err, "query metadata") + } + return metas, nil +} + +func (spec *SSplitTableSpec) GetTableSpec(meta STableMetadata) *sqlchemy.STableSpec { + tbSpec := *spec.tableSpec + return tbSpec.Clone(meta.Table, meta.Start) +} diff --git a/pkg/util/splitable/purge.go b/pkg/util/splitable/purge.go new file mode 100644 index 0000000000..2771622f9b --- /dev/null +++ b/pkg/util/splitable/purge.go @@ -0,0 +1,54 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package splitable + +import ( + "fmt" + "time" + + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/sqlchemy" +) + +func (t *SSplitTableSpec) Purge() error { + if t.maxSegments <= 0 { + return nil + } + metas, err := t.GetTableMetas() + if err != nil { + return errors.Wrap(err, "GetTableMetas") + } + if t.maxSegments >= len(metas) { + return nil + } + for i := 0; i < len(metas)-t.maxSegments; i += 1 { + dropSQL := fmt.Sprintf("DROP TABLE `%s`", metas[i].Table) + log.Infof("Ready to drop table: %s", dropSQL) + _, err := sqlchemy.Exec(dropSQL) + if err != nil { + return errors.Wrap(err, "sqlchemy.Exec") + } + _, err = t.metaSpec.Update(&metas[i], func() error { + metas[i].DeleteAt = time.Now() + metas[i].Deleted = true + return nil + }) + if err != nil { + return errors.Wrap(err, "metaSpec.Update") + } + } + return nil +} diff --git a/pkg/util/splitable/splitable.go b/pkg/util/splitable/splitable.go new file mode 100644 index 0000000000..63b6f6f9d4 --- /dev/null +++ b/pkg/util/splitable/splitable.go @@ -0,0 +1,152 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package splitable + +import ( + "database/sql" + "fmt" + "reflect" + "time" + + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/util/reflectutils" + "yunion.io/x/sqlchemy" +) + +type SSplitTableSpec struct { + indexField string + dateField string + tableName string + tableSpec *sqlchemy.STableSpec + metaSpec *sqlchemy.STableSpec + maxDuration time.Duration + maxSegments int +} + +func (t *SSplitTableSpec) DataType() reflect.Type { + return t.tableSpec.DataType() +} + +func (t *SSplitTableSpec) ColumnSpec(name string) sqlchemy.IColumnSpec { + return t.tableSpec.ColumnSpec(name) +} + +func (t *SSplitTableSpec) Name() string { + return t.tableName +} + +func (t *SSplitTableSpec) Columns() []sqlchemy.IColumnSpec { + return t.tableSpec.Columns() +} + +func (t *SSplitTableSpec) PrimaryColumns() []sqlchemy.IColumnSpec { + return t.tableSpec.PrimaryColumns() +} + +func (t *SSplitTableSpec) Expression() string { + metas, err := t.GetTableMetas() + if err != nil { + return fmt.Sprintf("`%s`", t.tableName) + } + tss := make([]sqlchemy.IQuery, 0) + for _, meta := range metas { + ts := t.GetTableSpec(meta) + tss = append(tss, ts.Query()) + } + union, err := sqlchemy.UnionWithError(tss...) + if err != nil { + return fmt.Sprintf("`%s`", t.tableName) + } + return union.Expression() +} + +func (t *SSplitTableSpec) Instance() *sqlchemy.STable { + return sqlchemy.NewTableInstance(t) +} + +func (t *SSplitTableSpec) DropForeignKeySQL() []string { + return t.tableSpec.DropForeignKeySQL() +} + +func (t *SSplitTableSpec) AddIndex(unique bool, cols ...string) bool { + metas, err := t.GetTableMetas() + if err != nil { + return false + } + var ret bool + for _, meta := range metas { + ts := t.GetTableSpec(meta) + if !ts.AddIndex(unique, cols...) { + ret = false + break + } + } + return ret +} + +func (t *SSplitTableSpec) Fetch(dt interface{}) error { + vs := reflectutils.FetchStructFieldValueSet(reflect.Indirect(reflect.ValueOf(dt))) + idxVal, ok := vs.GetValue(t.indexField) + if !ok { + return errors.Wrap(errors.ErrNotFound, "GetValue") + } + idxInt := idxVal.Int() + metas, err := t.GetTableMetas() + if err != nil { + return errors.Wrap(err, "GetTableMetas") + } + for _, meta := range metas { + if idxInt >= meta.Start && (meta.End == 0 || meta.End >= idxInt) { + ts := t.GetTableSpec(meta) + return ts.Fetch(dt) + } + } + return sql.ErrNoRows +} + +func NewSplitTableSpec(s interface{}, name string, indexField string, dateField string, maxDuration time.Duration, maxSegments int) (*SSplitTableSpec, error) { + spec := sqlchemy.NewTableSpecFromStruct(s, name) + indexCol := spec.ColumnSpec(indexField) + if indexCol == nil { + return nil, errors.Wrapf(errors.ErrNotFound, "indexField %s not found", indexField) + } + if !indexCol.IsPrimary() { + return nil, errors.Wrapf(errors.ErrInvalidStatus, "indexField %s not primary", indexField) + } + if intCol, ok := indexCol.(*sqlchemy.SIntegerColumn); !ok { + return nil, errors.Wrapf(errors.ErrInvalidStatus, "indexField %s not integer", indexField) + } else if !intCol.IsAutoIncrement { + return nil, errors.Wrapf(errors.ErrInvalidStatus, "indexField %s not auto_increment", indexField) + } + dateCol := spec.ColumnSpec(dateField) + if dateCol == nil { + return nil, errors.Wrapf(errors.ErrNotFound, "dateField %s not found", dateField) + } + if _, ok := dateCol.(*sqlchemy.SDateTimeColumn); !ok { + return nil, errors.Wrapf(errors.ErrInvalidStatus, "dateField %s not datetime column", dateField) + } + + metaSpec := sqlchemy.NewTableSpecFromStruct(&STableMetadata{}, fmt.Sprintf("%s_metadata", name)) + + return &SSplitTableSpec{ + indexField: indexField, + dateField: dateField, + tableName: name, + tableSpec: spec, + metaSpec: metaSpec, + maxDuration: maxDuration, + maxSegments: maxSegments, + }, nil +} diff --git a/pkg/util/splitable/sync.go b/pkg/util/splitable/sync.go new file mode 100644 index 0000000000..735a634aca --- /dev/null +++ b/pkg/util/splitable/sync.go @@ -0,0 +1,135 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package splitable + +import ( + "fmt" + "time" + + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/util/timeutils" + "yunion.io/x/sqlchemy" +) + +func (spec *SSplitTableSpec) Sync() error { + err := spec.metaSpec.Sync() + if err != nil { + return errors.Wrap(err, "metaSpec.Sync") + } + metas, err := spec.GetTableMetas() + if err != nil { + return errors.Wrap(err, "GetTableMetas") + } + if len(metas) == 0 { + // init the first metadata record + fakeMeta := STableMetadata{ + Table: spec.tableName, + } + tbl := spec.GetTableSpec(fakeMeta) + if tbl.Exists() { + err := tbl.Sync() + if err != nil { + return errors.Wrap(err, "Sync") + } + var minIndex int64 + var minDate time.Time + ti := tbl.Instance() + q := ti.Query(sqlchemy.MIN("min_index", ti.Field(spec.indexField)), sqlchemy.MIN("min_date", ti.Field(spec.dateField))) + r := q.Row() + err = r.Scan(&minIndex, &minDate) + if err != nil { + return errors.Wrap(err, "minIndex minDate") + } + fakeMeta.Start = minIndex + fakeMeta.StartDate = minDate + err = spec.metaSpec.Insert(&fakeMeta) + if err != nil { + return errors.Wrap(err, "insert init metadata") + } + } + } else { + for i := range metas { + subSpec := spec.GetTableSpec(metas[i]) + err := subSpec.Sync() + if err != nil { + return errors.Wrap(err, "Sync") + } + } + } + return nil +} + +func (spec *SSplitTableSpec) CheckSync() error { + err := spec.metaSpec.CheckSync() + if err != nil { + return errors.Wrap(err, "metaSpec.CheckSync") + } + metas, err := spec.GetTableMetas() + if err != nil { + return errors.Wrap(err, "GetTableMetas") + } + for i := range metas { + subSpec := spec.GetTableSpec(metas[i]) + err := subSpec.CheckSync() + if err != nil { + return errors.Wrap(err, "GetTableSpec") + } + } + return nil +} + +func (spec *SSplitTableSpec) SyncSQL() []string { + sqls := spec.metaSpec.SyncSQL() + if spec.metaSpec.Exists() { + metas, err := spec.GetTableMetas() + if err != nil { + log.Errorf("GetTableMetas fail %s", err) + return nil + } else if len(metas) > 0 { + for i := range metas { + subSpec := spec.GetTableSpec(metas[i]) + nsql := subSpec.SyncSQL() + sqls = append(sqls, nsql...) + } + return sqls + } + } + + fakeMeta := STableMetadata{ + Table: spec.tableName, + } + tbl := spec.GetTableSpec(fakeMeta) + if tbl.Exists() { + nsql := tbl.SyncSQL() + if len(nsql) > 0 { + sqls = append(sqls, nsql...) + } + var minIndex int64 + var minDate time.Time + ti := tbl.Instance() + q := ti.Query(sqlchemy.MIN("min_index", ti.Field(spec.indexField)), sqlchemy.MIN("min_date", ti.Field(spec.dateField))) + r := q.Row() + err := r.Scan(&minIndex, &minDate) + if err != nil { + log.Errorf("query minIndex minDate fail %s", err) + } else { + minDateStr := timeutils.MysqlTime(minDate) + sql := fmt.Sprintf("INSERT INTO `%s`(`table`, `start`, `start_date`, `deleted`, `created_at`) VALUES('%s', %d, '%s', 0, '%s')", spec.metaSpec.Name(), spec.tableName, minIndex, minDateStr, minDateStr) + sqls = append(sqls, sql) + } + } + return sqls +} diff --git a/pkg/util/splitable/update.go b/pkg/util/splitable/update.go new file mode 100644 index 0000000000..8a53092bc2 --- /dev/null +++ b/pkg/util/splitable/update.go @@ -0,0 +1,36 @@ +// Copyright 2019 Yunion +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package splitable + +import ( + "yunion.io/x/pkg/errors" + "yunion.io/x/sqlchemy" +) + +func (t *SSplitTableSpec) InsertOrUpdate(dt interface{}) error { + return errors.ErrNotSupported +} + +func (t *SSplitTableSpec) Update(dt interface{}, onUpdate func() error) (sqlchemy.UpdateDiffs, error) { + return nil, errors.ErrNotSupported +} + +func (t *SSplitTableSpec) Increment(diff, target interface{}) error { + return errors.ErrNotSupported +} + +func (t *SSplitTableSpec) Decrement(diff, target interface{}) error { + return errors.ErrNotSupported +}