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, )
| 191 | |
| 192 | // Helper function to build the Arrow record and register the view |
| 193 | func (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) |
no test coverage detected