forked from dolthub/go-mysql-server
-
Notifications
You must be signed in to change notification settings - Fork 0
/
window_iter.go
150 lines (134 loc) · 3.49 KB
/
window_iter.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
// Copyright 2022 Dolthub, Inc.
//
// 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 aggregation
import (
"errors"
"io"
"github.com/gabereiser/go-mysql-server/sql"
)
// WindowIter is a wrapper that evaluates a set of WindowPartitionIter.
//
// The current implementation has 3 steps:
// 1. Materialize [iter] and duplicate a sql.WindowBuffer for each partition.
// 2. Collect rows from child partitions.
// 3. Rearrange partition results into the projected ordering given by [outputOrdinals].
//
// We assume [outputOrdinals] is appropriately sized for [partitionIters].
type WindowIter struct {
partitionIters []*WindowPartitionIter
outputOrdinals [][]int
iter sql.RowIter
initialized bool
}
func NewWindowIter(partitionIters []*WindowPartitionIter, outputOrdinals [][]int, iter sql.RowIter) *WindowIter {
return &WindowIter{
partitionIters: partitionIters,
outputOrdinals: outputOrdinals,
iter: iter,
}
}
var _ sql.RowIter = (*WindowIter)(nil)
var _ sql.Disposable = (*WindowIter)(nil)
// Close implements sql.RowIter
func (i *WindowIter) Close(ctx *sql.Context) error {
i.Dispose()
var err error
for _, p := range i.partitionIters {
e := p.Close(ctx)
if err == nil && e != nil {
err = e
}
}
return err
}
// Dispose implements sql.Disposable
func (i *WindowIter) Dispose() {
for _, p := range i.partitionIters {
p.Dispose()
}
return
}
// Next implements sql.RowIter
func (i *WindowIter) Next(ctx *sql.Context) (sql.Row, error) {
if !i.initialized {
err := i.initializeIters(ctx)
if err != nil {
return nil, err
}
}
row := make(sql.Row, i.size())
for j, pIter := range i.partitionIters {
res, err := pIter.Next(ctx)
if err != nil {
return nil, err
}
for k, idx := range i.outputOrdinals[j] {
row[idx] = res[k]
}
}
return row, nil
}
func (i *WindowIter) size() int {
size := -1
for _, i := range i.outputOrdinals {
for _, j := range i {
if j > size {
size = j
}
}
}
return size + 1
}
// initializeIters materializes and copies the input buffer into each
// WindowPartitionIter.
// TODO: share the child buffer and sort/partition inbetween WindowPartitionIters
func (i *WindowIter) initializeIters(ctx *sql.Context) error {
buf := make(sql.WindowBuffer, 0)
var row sql.Row
var err error
for {
// drain child iter into reusable buffer
row, err = i.iter.Next(ctx)
if errors.Is(err, io.EOF) {
break
}
if err != nil {
return err
}
buf = append(buf, row)
}
for _, i := range i.partitionIters {
// each iter has its own copy of input buffer
i.child = &windowBufferIter{buf: buf}
}
i.initialized = true
return nil
}
// windowBufferIter bridges an in-memory buffer to the sql.RowIter interface
type windowBufferIter struct {
buf sql.WindowBuffer
pos int
}
func (i *windowBufferIter) Next(ctx *sql.Context) (sql.Row, error) {
if i.pos >= len(i.buf) {
return nil, io.EOF
}
row := i.buf[i.pos]
i.pos++
return row, nil
}
func (i *windowBufferIter) Close(ctx *sql.Context) error {
i.buf = nil
return nil
}