diff --git a/executor.go b/executor.go index 0f52e09d8..ade3b9001 100644 --- a/executor.go +++ b/executor.go @@ -2182,16 +2182,14 @@ func (e *executor) executeTopKShardTime(ctx context.Context, tx Tx, filter *Row, return nil, newNotFoundError(ErrFieldNotFound, field) } - // Check the time quantum. - quantum := f.TimeQuantum() - if quantum == "" { - // ???????? - return nil, nil + views, err := f.viewsByTimeRange(from, to) + if err != nil { + return nil, err } // Fetch fragments. var fragments []*fragment - for _, view := range viewsByTimeRange(viewStandard, from, to, quantum) { + for _, view := range views { f := e.Holder.fragment(index, field, view, shard) if f == nil { continue @@ -3746,48 +3744,9 @@ func (e *executor) executeRowsShard(ctx context.Context, qcx *Qcx, index string, return nil, errors.Wrap(err, "parsing to time") } } - - // Calculate the views for a range as long as some piece of the range - // (from/to) are specified, or if there's no standard view to represent - // all dates. - if !fromTime.IsZero() || !toTime.IsZero() || f.options.NoStandardView { - // If no quantum exists then return an empty result set. - q := f.TimeQuantum() - if q == "" { - return rowIDs, nil - } - - // Get min/max based on existing views. - var vs []string - for _, v := range f.views() { - vs = append(vs, v.name) - } - min, max := minMaxViews(vs, q) - - // If min/max are empty, there were no time views. - if min == "" || max == "" { - return rowIDs, nil - } - - // Convert min/max from string to time.Time. - minTime, err := timeOfView(min, false) - if err != nil { - return rowIDs, errors.Wrapf(err, "getting min time from view: %s", min) - } - if fromTime.IsZero() || fromTime.Before(minTime) { - fromTime = minTime - } - - maxTime, err := timeOfView(max, true) - if err != nil { - return rowIDs, errors.Wrapf(err, "getting max time from view: %s", max) - } - if toTime.IsZero() || toTime.After(maxTime) { - toTime = maxTime - } - - // Determine the views based on the specified time range. - views = viewsByTimeRange(viewStandard, fromTime, toTime, q) + views, err = f.viewsByTimeRange(fromTime, toTime) + if err != nil { + return nil, err } default: return nil, errors.Errorf("%s fields not supported by Rows() query", f.Type()) @@ -4585,21 +4544,12 @@ func (e *executor) executeRowShard(ctx context.Context, qcx *Qcx, index string, return row, err } - // If no quantum exists then return an empty bitmap. - q := f.TimeQuantum() - if q == "" { - return &Row{}, nil - } - - // Set maximum "to" value if only "from" is set. We don't need to worry - // about setting the minimum "from" since it is the zero value if omitted. - if toTime.IsZero() { - // Set the end timestamp to current time + 1 day, in order to account for timezone differences. - toTime = time.Now().AddDate(0, 0, 1) + views, err := f.viewsByTimeRange(fromTime, toTime) + if err != nil { + return nil, err } // Union bitmaps across all time-based views. - views := viewsByTimeRange(viewStandard, fromTime, toTime, q) rows := make([]*Row, 0, len(views)) tx, finisher, err := qcx.GetTx(Txo{Write: !writable, Index: idx, Shard: shard}) defer finisher(&err0) @@ -7856,14 +7806,14 @@ func newGroupByIterator(executor *executor, qcx *Qcx, rowIDs []RowIDs, children idx := holder.Index(index) var ( - fieldName string - viewName string - ok bool - views []string - isTimeField bool + fieldName string + viewName string + ok bool + views []string ) ignorePrev := false for i, call := range children { + var isTimeField bool if fieldName, ok = call.Args["_field"].(string); !ok { return nil, errors.Errorf("%s call must have field with valid (string) field name. Got %v of type %[2]T", call.Name, call.Args["_field"]) } @@ -7907,7 +7857,12 @@ func newGroupByIterator(executor *executor, qcx *Qcx, rowIDs []RowIDs, children } if hasTo || hasFrom { - views = viewsByTimeRange(viewStandard, fromTime, toTime, field.TimeQuantum()) + // Determine the views based on the specified time range. + var err error + views, err = field.viewsByTimeRange(fromTime, toTime) + if err != nil { + return nil, err + } isTimeField = true } else { viewName = viewStandard diff --git a/executor_test.go b/executor_test.go index 6f0e8839d..069f27b7b 100644 --- a/executor_test.go +++ b/executor_test.go @@ -469,21 +469,17 @@ func TestExecutor(t *testing.T) { t.Run("Range", func(t *testing.T) { t.Run("RowIDColumnID", func(t *testing.T) { - // Create a timestamp just out of the current date + 1 day timestamp (default end timestamp). - nextDayExclusive := time.Now().AddDate(0, 0, 2) - - writeQuery := fmt.Sprintf(` + writeQuery := ` Set(2, f=1, 1999-12-31T00:00) Set(3, f=1, 2000-01-01T00:00) Set(4, f=1, 2000-01-02T00:00) Set(5, f=1, 2000-02-01T00:00) Set(6, f=1, 2001-01-01T00:00) Set(7, f=1, 2002-01-01T02:00) - Set(8, f=1, %s) Set(2, f=1, 1999-12-30T00:00) Set(2, f=1, 2002-02-01T00:00) - Set(2, f=10, 2001-01-01T00:00)`, nextDayExclusive.Format("2006-01-02T15:04")) + Set(2, f=10, 2001-01-01T00:00)` readQueries := []string{ `Row(f=1, from=1999-12-31T00:00, to=2002-01-01T03:00)`, `Row(f=1, from=1999-12-31T00:00)`, @@ -5766,6 +5762,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) { defer c.Close() c.CreateField(t, "i", pilosa.IndexOptions{}, "general") c.CreateField(t, "i", pilosa.IndexOptions{}, "sub") + c.CreateField(t, "i", pilosa.IndexOptions{}, "tq", pilosa.OptFieldTypeTime("YMDH", "0")) c.CreateField(t, "i", pilosa.IndexOptions{}, "v", pilosa.OptFieldTypeInt(0, 1000)) c.ImportBits(t, "i", "general", [][2]uint64{ {10, 0}, @@ -5786,6 +5783,7 @@ func TestExecutor_Execute_GroupBy(t *testing.T) { {110, 2}, {110, 0}, }) + if _, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(0, v=10)`}); err != nil { t.Fatal(err) } else if _, err := c.GetNode(0).API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(1, v=100)`}); err != nil { @@ -6198,6 +6196,38 @@ func TestExecutor_Execute_GroupBy(t *testing.T) { ) }) + // Create some time-quantum data: + c.Query(t, "i", "Set(0, tq=1, 2022-01-01T01:01)") + c.Query(t, "i", "Set(1, tq=1, 2021-01-01T01:01)") + t.Run("GroupByWithTime", func(t *testing.T) { + expected := map[string][]pilosa.GroupCount{ + // no time specified + "GroupBy(Rows(tq), Rows(general))": { + {Group: []pilosa.FieldRow{{Field: "tq", RowID: 1}, {Field: "general", RowID: 10}}, Count: 2}, + }, + // time specified but includes all data + "GroupBy(Rows(tq, from=2020-01-01T01:01), Rows(general))": { + {Group: []pilosa.FieldRow{{Field: "tq", RowID: 1}, {Field: "general", RowID: 10}}, Count: 2}, + }, + // same but in a different order + "GroupBy(Rows(general), Rows(tq, from=2020-01-01T01:01))": { + {Group: []pilosa.FieldRow{{Field: "general", RowID: 10}, {Field: "tq", RowID: 1}}, Count: 2}, + }, + // time excludes any data + "GroupBy(Rows(general), Rows(tq, from=2022-01-01T01:01))": { + {Group: []pilosa.FieldRow{{Field: "general", RowID: 10}, {Field: "tq", RowID: 1}}, Count: 1}, + }, + // limit excludes all data + "GroupBy(Rows(general), Rows(tq, from=2023-01-01T01:01))": {}, + } + + for query, want := range expected { + results := c.Query(t, "i", query).Results[0].(*pilosa.GroupCounts).Groups() + t.Logf("query %q", query) + test.CheckGroupBy(t, want, results) + } + + }) } for _, size := range []int{1, 3} { t.Run(fmt.Sprintf("%d_nodes", size), func(t *testing.T) { diff --git a/field.go b/field.go index 8c404068c..df87eea91 100644 --- a/field.go +++ b/field.go @@ -968,6 +968,55 @@ func (f *Field) TimeQuantum() TimeQuantum { return f.options.TimeQuantum } +// viewsByTimeRange is a wrapper on the non-method viewsByTimeRange, which +// computes views for a specific field for a given time range. The difference +// is that, as a Field operation, it can return "standard" for a view that +// covers the whole time range, if the field supports a standard view, and +// can automatically coerce from/to times to match the actual range present +// in the field. +func (f *Field) viewsByTimeRange(from, to time.Time) (views []string, err error) { + // If we can't find time views at all, we'll yield "standard" regardless. + // It's the least-bad answer, I think. + q := f.TimeQuantum() + if q == "" { + return []string{viewStandard}, nil + } + + // Get min/max based on existing views. + var vs []string + for _, v := range f.views() { + vs = append(vs, v.name) + } + min, max := minMaxViews(vs, q) + + // If min/max are empty, there were no time views. + if min == "" || max == "" { + return []string{viewStandard}, nil + } + + wasZero := from.IsZero() && to.IsZero() + // Convert min/max from string to time.Time. + minTime, err := timeOfView(min, false) + if err != nil { + return nil, errors.Wrapf(err, "getting min time from view: %s", min) + } + if from.IsZero() || from.Before(minTime) { + from = minTime + } + + maxTime, err := timeOfView(max, true) + if err != nil { + return nil, errors.Wrapf(err, "getting max time from view: %s", max) + } + if to.IsZero() || to.After(maxTime) { + to = maxTime + } + if (wasZero || (from == minTime && to == maxTime)) && !f.Options().NoStandardView { + return []string{viewStandard}, nil + } + return viewsByTimeRange(viewStandard, from, to, q), nil +} + // RowTime gets the row at the particular time with the granularity specified by // the quantum. func (f *Field) RowTime(tx Tx, rowID uint64, time time.Time, quantum string) (*Row, error) { diff --git a/field_internal_test.go b/field_internal_test.go index fb2f01060..268ca95d6 100644 --- a/field_internal_test.go +++ b/field_internal_test.go @@ -911,3 +911,70 @@ func TestField_SaveMeta(t *testing.T) { t.Fatalf("expected value after reopen to be: %d, got: %d", val, rslt) } } + +func TestFieldViewsByTimeRange(t *testing.T) { + f := OpenField(t, OptFieldTypeTime("YMD", "0", false)) + for _, date := range []string{ + // a handful of YMD parameters describing dates that we could have data for + "2021", + "202112", + "20211229", + "20211230", + "20211231", + "2022", + "202201", + "20220101", + "20220102", + } { + _, err := f.createViewIfNotExists(viewStandard + "_" + date) + if err != nil { + t.Fatalf("creating view for %s: %v", date, err) + } + } + var testCases = []struct { + from, to string + expected []string + }{ + {"", "", []string{"standard"}}, + {"2020-12-31T00:00", "2023-01-03T00:00", []string{"standard"}}, + {"2021-01-01T00:00", "2022-01-01T00:00", []string{"standard_2021"}}, + {"2021-01-01T00:00", "2022-01-02T00:00", []string{"standard_2021", "standard_20220101"}}, + {"", "2022-01-02T00:00", []string{"standard_2021", "standard_20220101"}}, + {"2021-12-01T00:00", "", []string{"standard_202112", "standard_2022"}}, + {"2021-12-30T00:00", "2022-02-01T00:00", []string{"standard_20211230", "standard_20211231", "standard_202201"}}, + } + for _, tc := range testCases { + t.Logf("checking %q to %q", tc.from, tc.to) + var fromTime, toTime time.Time + var err error + if tc.from != "" { + fromTime, err = time.Parse("2006-01-02T15:04", tc.from) + if err != nil { + t.Fatalf("invalid time %q: %v", tc.from, err) + } + } + if tc.to != "" { + toTime, err = time.Parse("2006-01-02T15:04", tc.to) + if err != nil { + t.Fatalf("invalid time %q: %v", tc.to, err) + } + } + views, err := f.viewsByTimeRange(fromTime, toTime) + if err != nil { + t.Fatalf("unexpected error getting views for %s-%s: %v", tc.from, tc.to, err) + } + for i, v := range tc.expected { + if len(views) <= i { + t.Fatalf("expected view %q, didn't get it", v) + } else { + if views[i] != v { + t.Fatalf("expected view %q, got %q", v, views[i]) + } + } + } + if len(views) > len(tc.expected) { + t.Fatalf("unexpected view %q", views[len(tc.expected)]) + } + t.Logf("views: %v", views) + } +}