MCPcopy Create free account
hub / github.com/apecloud/myduckserver / prepareArrowView

Method prepareArrowView

delta/controller.go:193–264  ·  view source on GitHub ↗

Helper function to build the Arrow record and register the view

(
	ctx *sql.Context,
	conn *stdsql.Conn,
	table tableIdentifier,
	appender *DeltaAppender,
	fieldOffset int,
	fieldIndices []int,
)

Source from the content-addressed store, hash-verified

191
192// Helper function to build the Arrow record and register the view
193func (c *DeltaController) prepareArrowView(
194 ctx *sql.Context,
195 conn *stdsql.Conn,
196 table tableIdentifier,
197 appender *DeltaAppender,
198 fieldOffset int,
199 fieldIndices []int,
200) (viewName string, close func(), err error) {
201 record := appender.Build()
202
203 // fmt.Println("record:", record)
204
205 var ar *duckdb.Arrow
206 err = conn.Raw(func(driverConn any) error {
207 var err error
208 ar, err = duckdb.NewArrowFromConn(driverConn.(*duckdb.Conn))
209 return err
210 })
211 if err != nil {
212 record.Release()
213 return "", nil, err
214 }
215
216 // Project the fields before registering the Arrow record into DuckDB.
217 // Currently, this is necessary because RegisterView uses `arrow_scan` instead of `arrow_scan_dumb` under the hood.
218 // The former allows projection & filter pushdown, but the implementation does not work as expected in some cases.
219 schema := record.Schema()
220 if fieldOffset > 0 {
221 fields := schema.Fields()[fieldOffset:]
222 schema = arrow.NewSchema(fields, nil)
223 columns := record.Columns()[fieldOffset:]
224 projected := array.NewRecord(schema, columns, record.NumRows())
225 record.Release()
226 record = projected
227 } else if len(fieldIndices) > 0 {
228 fields := make([]arrow.Field, len(fieldIndices))
229 columns := make([]arrow.Array, len(fieldIndices))
230 for i, idx := range fieldIndices {
231 fields[i] = schema.Field(idx)
232 columns[i] = record.Column(idx)
233 }
234 schema = arrow.NewSchema(fields, nil)
235 projected := array.NewRecord(schema, columns, record.NumRows())
236 record.Release()
237 record = projected
238 }
239
240 reader, err := array.NewRecordReader(schema, []arrow.Record{record})
241 if err != nil {
242 record.Release()
243 return "", nil, err
244 }
245
246 // Register the Arrow view
247 hash := maphash.String(c.seed, table.dbName+"\x00"+table.tableName)
248 viewName = "__sys_view_arrow_delta_" + strconv.FormatUint(hash, 16) + "__"
249
250 release, err := ar.RegisterView(reader, viewName)

Callers 4

handleInsertOnlyMethod · 0.95
handleDeleteOnlyMethod · 0.95
handleZeroDeleteMethod · 0.95

Calls 7

ExecContextMethod · 0.80
FieldsMethod · 0.65
FieldMethod · 0.65
BuildMethod · 0.45
ReleaseMethod · 0.45
SchemaMethod · 0.45
StringMethod · 0.45

Tested by

no test coverage detected