(object *data.Object, recs []super.Value)
| 128 | } |
| 129 | |
| 130 | func (w *Writer) writeObject(object *data.Object, recs []super.Value) error { |
| 131 | var zr sio.Reader |
| 132 | if w.inputSorted { |
| 133 | zr = sbuf.NewArray(recs) |
| 134 | } else { |
| 135 | done := make(chan struct{}) |
| 136 | go func() { |
| 137 | zr = w.comparator.SortStableReader(recs) |
| 138 | close(done) |
| 139 | }() |
| 140 | select { |
| 141 | case <-done: |
| 142 | case <-w.ctx.Done(): |
| 143 | return w.ctx.Err() |
| 144 | } |
| 145 | } |
| 146 | writer, err := object.NewWriter(w.ctx, w.pool.engine, w.pool.DataPath, w.pool.SortKeys.Primary(), w.pool.SeekStride) |
| 147 | if err != nil { |
| 148 | return err |
| 149 | } |
| 150 | if err := sio.CopyWithContext(w.ctx, writer, zr); err != nil { |
| 151 | writer.Abort() |
| 152 | return err |
| 153 | } |
| 154 | if err := writer.Close(w.ctx); err != nil { |
| 155 | return err |
| 156 | } |
| 157 | w.stats.Accumulate(ImportStats{ |
| 158 | ObjectsWritten: 1, |
| 159 | RecordBytesWritten: writer.BytesWritten(), |
| 160 | RecordsWritten: int64(writer.RecordsWritten()), |
| 161 | }) |
| 162 | return nil |
| 163 | } |
| 164 | |
| 165 | func (w *Writer) Stats() ImportStats { |
| 166 | return w.stats.Copy() |
no test coverage detected