-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathpool.go
More file actions
72 lines (65 loc) · 1.77 KB
/
Copy pathpool.go
File metadata and controls
72 lines (65 loc) · 1.77 KB
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
package roost
import (
"io"
"github.com/apache/arrow-go/v18/arrow"
)
// objState is the per-object encode state shared by all of an object's ops. All
// ops for one object are routed to a single worker (affinity), so failed is only
// ever touched by that one goroutine - no synchronization needed.
type objState struct {
obj ObjectEncoder
wc io.WriteCloser
cw *countWriter
failed bool
}
// encodeOp is one unit of work for an encode worker: write rec as a row group,
// or, when rec is nil, finalize the object (footer + upload).
type encodeOp struct {
os *objState
rec arrow.RecordBatch
}
// encodeWorker drains one worker channel. A fixed set of these (sized by
// WithEncodeConcurrency) does all compression and upload off the Writer's lock,
// so the goroutine count is bounded regardless of how many objects churn. Each
// object is pinned to one worker, so its row groups are written in order.
func (w *Writer[T]) encodeWorker(ch <-chan encodeOp) {
defer w.wg.Done()
for op := range ch {
os := op.os
if op.rec != nil {
if !os.failed {
if err := os.obj.Write(op.rec); err != nil {
os.failed = true
w.setEncErr(err)
}
}
op.rec.Release()
continue
}
// rec == nil: finalize.
if err := os.obj.Close(); err != nil && !os.failed {
os.failed = true
w.setEncErr(err)
}
if !os.failed {
w.stats.addObject(os.cw.n)
}
if err := os.wc.Close(); err != nil {
w.setEncErr(err)
}
}
}
// setEncErr records the first error seen by any encode worker.
func (w *Writer[T]) setEncErr(err error) {
w.encMu.Lock()
if w.encErr == nil {
w.encErr = err
}
w.encMu.Unlock()
}
// encErrLoad returns the first encode error seen so far, if any.
func (w *Writer[T]) encErrLoad() error {
w.encMu.Lock()
defer w.encMu.Unlock()
return w.encErr
}