diff --git a/executor.go b/executor.go index 722f63d32..ff5317b04 100644 --- a/executor.go +++ b/executor.go @@ -7262,9 +7262,11 @@ func newGroupByIterator(executor *executor, qcx *Qcx, rowIDs []RowIDs, children idx := holder.Index(index) var ( - fieldName string - viewName string - ok bool + fieldName string + viewName string + ok bool + views []string + isTimeField bool ) ignorePrev := false for i, call := range children { @@ -7278,9 +7280,42 @@ func newGroupByIterator(executor *executor, qcx *Qcx, rowIDs []RowIDs, children gbi.fields[i].Field = fieldName switch field.Type() { - case FieldTypeSet, FieldTypeTime, FieldTypeMutex, FieldTypeBool: + case FieldTypeSet, FieldTypeMutex, FieldTypeBool: viewName = viewStandard + case FieldTypeTime: + var ( + err error + v interface{} + ) + // Parse "from" time, if set. + var ( + hasFrom bool + fromTime time.Time + ) + if v, hasFrom = call.Args["from"]; hasFrom { + if fromTime, err = parseTime(v); err != nil { + return nil, errors.Wrap(err, "parsing from time") + } + } + + // Parse "to" time, if set. + var ( + hasTo bool + toTime time.Time + ) + if v, hasTo = call.Args["to"]; hasTo { + if toTime, err = parseTime(v); err != nil { + return nil, errors.Wrap(err, "parsing to time") + } + } + + if hasTo || hasFrom { + views = viewsByTimeRange(viewStandard, fromTime, toTime, field.TimeQuantum()) + isTimeField = true + } else { + viewName = viewStandard + } case FieldTypeInt: viewName = viewBSIGroupPrefix + fieldName @@ -7289,11 +7324,6 @@ func newGroupByIterator(executor *executor, qcx *Qcx, rowIDs []RowIDs, children call.Name, strings.Join([]string{FieldTypeSet, FieldTypeTime, FieldTypeMutex, FieldTypeBool, FieldTypeInt}, ",")) } - // Fetch fragment. - frag := holder.fragment(index, fieldName, viewName, shard) - if frag == nil { // this means this whole shard doesn't have all it needs to continue - return nil, nil - } filters := []roaring.BitmapFilter{} if len(rowIDs[i]) > 0 { filters = append(filters, roaring.NewBitmapRowsFilter(rowIDs[i])) @@ -7305,9 +7335,35 @@ func newGroupByIterator(executor *executor, qcx *Qcx, rowIDs []RowIDs, children } defer finisher(&err0) - gbi.rowIters[i], err = frag.rowIterator(tx, i != 0, filters...) - if err != nil { - return nil, err + // Fetch fragment(s), get rowIterator + if isTimeField { + var fragments []*fragment + for _, viewName := range views { + fragment := holder.fragment(index, fieldName, viewName, shard) + if fragment != nil { + fragments = append(fragments, fragment) + } + } + if len(fragments) == 0 { + // whole shard doesn't have all it needs to continue ? + return nil, nil + } + + gbi.rowIters[i], err = timeFragmentsRowIterator(fragments, tx, i != 0, filters...) + if err != nil { + return nil, err + } + } else { + frag := holder.fragment(index, fieldName, viewName, shard) + if frag == nil { // this means this whole shard doesn't have all it needs to continue + return nil, nil + } + + gbi.rowIters[i], err = frag.rowIterator(tx, i != 0, filters...) + if err != nil { + return nil, err + } + } prev, hasPrev, err := call.UintArg("previous") diff --git a/fragment.go b/fragment.go index aa93dcc9a..031043390 100644 --- a/fragment.go +++ b/fragment.go @@ -3329,6 +3329,75 @@ func (f *fragment) rowIterator(tx Tx, wrap bool, filters ...roaring.BitmapFilter return f.setRowIterator(tx, wrap, filters...) } +type timeRowIterator struct { + tx Tx + fragments []*fragment + rows []*Row + rowIDs [][]uint64 + cur int + wrap bool +} + +func timeFragmentsRowIterator(fragments []*fragment, tx Tx, wrap bool, filters ...roaring.BitmapFilter) (rowIterator, error) { + if len(fragments) == 0 { + return nil, fmt.Errorf("there should be at least 1 fragment") + } else if len(fragments) == 1 { + return fragments[0].setRowIterator(tx, wrap, filters...) + } + + it := &timeRowIterator{ + tx: tx, + fragments: fragments, + rows: make([]*Row, len(fragments)), + rowIDs: make([][]uint64, len(fragments)), + } + + for i, f := range fragments { + rows, err := f.rows(context.Background(), tx, 0, filters...) + if err != nil { + return nil, err + } + it.rowIDs[i] = rows + } + + return it, nil +} + +func (it *timeRowIterator) Seek(rowID uint64) { + rowIDs := it.rowIDs[0] + idx := sort.Search(len(rowIDs), func(i int) bool { + return rowIDs[i] >= rowID + }) + it.cur = idx +} + +func (it *timeRowIterator) Next() (r *Row, rowID uint64, _ *int64, wrapped bool, err error) { + rowIDs := it.rowIDs[0] + if it.cur >= len(rowIDs) { + if !it.wrap || len(rowIDs) == 0 { + return nil, 0, nil, true, nil + } + it.Seek(0) + wrapped = true + } + + id := rowIDs[it.cur] + // gather rows + for i, fragment := range it.fragments { + row, err := fragment.row(it.tx, id) + if err != nil { + return row, rowID, nil, wrapped, err + } + it.rows[i] = row + } + + // union rows + r = it.rows[0].Union(it.rows[1:]...) + + it.cur++ + return r, rowID, nil, wrapped, nil +} + type intRowIterator struct { f *fragment values int64Slice // sorted slice of int values