mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge pull request #9722 from swordqiu/automated-cherry-pick-of-#9719-upstream-release-3.7
Automated cherry pick of #9719: fix(cloudcommon): splitable may not initialize underlying table
This commit is contained in:
@@ -44,7 +44,7 @@ func (t *SSplitTableSpec) Insert(dt interface{}) error {
|
||||
newMeta := false
|
||||
if len(metas) > 0 {
|
||||
lastMeta := metas[len(metas)-1]
|
||||
if lastDate.Sub(lastMeta.StartDate) > t.maxDuration {
|
||||
if !lastMeta.StartDate.IsZero() && 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)))
|
||||
@@ -64,29 +64,49 @@ func (t *SSplitTableSpec) Insert(dt interface{}) error {
|
||||
}
|
||||
newMeta = true
|
||||
} else {
|
||||
if lastMeta.StartDate.IsZero() {
|
||||
indexCol := t.tableSpec.ColumnSpec(t.indexField)
|
||||
_, err = t.metaSpec.Update(&lastMeta, func() error {
|
||||
lastMeta.Start = indexCol.(*sqlchemy.SIntegerColumn).AutoIncrementOffset
|
||||
lastMeta.StartDate = lastDate
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "Update last meta")
|
||||
}
|
||||
}
|
||||
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)
|
||||
lastTableSpec, err = t.newTable(lastRecIndex, lastDate)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "insert new meta")
|
||||
return errors.Wrap(err, "newTable")
|
||||
}
|
||||
// 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)
|
||||
}
|
||||
|
||||
func (t *SSplitTableSpec) newTable(lastRecIndex int64, lastDate time.Time) (*sqlchemy.STableSpec, error) {
|
||||
// insert a new metadata
|
||||
meta := STableMetadata{
|
||||
Table: fmt.Sprintf("%s_%d", t.tableName, lastDate.Unix()),
|
||||
}
|
||||
if lastRecIndex > 0 {
|
||||
meta.Start = lastRecIndex + 1
|
||||
meta.StartDate = lastDate
|
||||
}
|
||||
err := t.metaSpec.Insert(&meta)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "insert new meta")
|
||||
}
|
||||
// create new table
|
||||
newTable := t.GetTableSpec(meta)
|
||||
err = newTable.Sync()
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "sync new table")
|
||||
}
|
||||
return newTable, nil
|
||||
}
|
||||
|
||||
@@ -140,7 +140,7 @@ func NewSplitTableSpec(s interface{}, name string, indexField string, dateField
|
||||
|
||||
metaSpec := sqlchemy.NewTableSpecFromStruct(&STableMetadata{}, fmt.Sprintf("%s_metadata", name))
|
||||
|
||||
return &SSplitTableSpec{
|
||||
sts := &SSplitTableSpec{
|
||||
indexField: indexField,
|
||||
dateField: dateField,
|
||||
tableName: name,
|
||||
@@ -148,5 +148,7 @@ func NewSplitTableSpec(s interface{}, name string, indexField string, dateField
|
||||
metaSpec: metaSpec,
|
||||
maxDuration: maxDuration,
|
||||
maxSegments: maxSegments,
|
||||
}, nil
|
||||
}
|
||||
|
||||
return sts, nil
|
||||
}
|
||||
|
||||
@@ -59,6 +59,11 @@ func (spec *SSplitTableSpec) Sync() error {
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "insert init metadata")
|
||||
}
|
||||
} else {
|
||||
_, err := spec.newTable(-1, time.Time{})
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "spec.newTable")
|
||||
}
|
||||
}
|
||||
} else {
|
||||
for i := range metas {
|
||||
@@ -81,11 +86,15 @@ func (spec *SSplitTableSpec) CheckSync() error {
|
||||
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")
|
||||
if len(metas) == 0 {
|
||||
return errors.Wrap(err, "empty metadata")
|
||||
} else {
|
||||
for i := range metas {
|
||||
subSpec := spec.GetTableSpec(metas[i])
|
||||
err := subSpec.CheckSync()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "GetTableSpec")
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
@@ -93,6 +102,8 @@ func (spec *SSplitTableSpec) CheckSync() error {
|
||||
|
||||
func (spec *SSplitTableSpec) SyncSQL() []string {
|
||||
sqls := spec.metaSpec.SyncSQL()
|
||||
zeroMeta := false
|
||||
|
||||
if spec.metaSpec.Exists() {
|
||||
metas, err := spec.GetTableMetas()
|
||||
if err != nil {
|
||||
@@ -105,7 +116,30 @@ func (spec *SSplitTableSpec) SyncSQL() []string {
|
||||
sqls = append(sqls, nsql...)
|
||||
}
|
||||
return sqls
|
||||
} else { // len(metas) == 0
|
||||
zeroMeta = true
|
||||
}
|
||||
} else {
|
||||
nsql := spec.metaSpec.SyncSQL()
|
||||
sqls = append(sqls, nsql...)
|
||||
zeroMeta = true
|
||||
}
|
||||
|
||||
if zeroMeta {
|
||||
indexCol := spec.tableSpec.ColumnSpec(spec.indexField)
|
||||
now := time.Now()
|
||||
meta := STableMetadata{
|
||||
Table: fmt.Sprintf("%s_%d", spec.tableName, now.Unix()),
|
||||
Start: indexCol.(*sqlchemy.SIntegerColumn).AutoIncrementOffset,
|
||||
}
|
||||
// insert the first meta
|
||||
sql := fmt.Sprintf("INSERT INTO `%s`(`table`, `deleted`, `created_at`) VALUES('%s', 0, '%s')", spec.metaSpec.Name(), meta.Table, timeutils.MysqlTime(now))
|
||||
sqls = append(sqls, sql)
|
||||
// create the first table
|
||||
newtable := spec.GetTableSpec(meta)
|
||||
nsql := newtable.SyncSQL()
|
||||
sqls = append(sqls, nsql...)
|
||||
return sqls
|
||||
}
|
||||
|
||||
fakeMeta := STableMetadata{
|
||||
|
||||
Reference in New Issue
Block a user