diff --git a/pkg/util/splitable/insert.go b/pkg/util/splitable/insert.go index 58bd0196fc..87d2bd3c76 100644 --- a/pkg/util/splitable/insert.go +++ b/pkg/util/splitable/insert.go @@ -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 +} diff --git a/pkg/util/splitable/splitable.go b/pkg/util/splitable/splitable.go index 63b6f6f9d4..9bc902066c 100644 --- a/pkg/util/splitable/splitable.go +++ b/pkg/util/splitable/splitable.go @@ -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 } diff --git a/pkg/util/splitable/sync.go b/pkg/util/splitable/sync.go index 735a634aca..815ae42faa 100644 --- a/pkg/util/splitable/sync.go +++ b/pkg/util/splitable/sync.go @@ -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{