| 23 | ) |
| 24 | |
| 25 | func (c *Client) Read(ctx context.Context, table *schema.Table, res chan<- arrow.RecordBatch) error { |
| 26 | colNames := make([]string, 0, len(table.Columns)) |
| 27 | for _, col := range table.Columns { |
| 28 | if c.pgType == pgTypeCrateDB { |
| 29 | colNames = append(colNames, pgx.Identifier{strings.Trim(col.Name, "_")}.Sanitize()) |
| 30 | continue |
| 31 | } |
| 32 | colNames = append(colNames, pgx.Identifier{col.Name}.Sanitize()) |
| 33 | } |
| 34 | cols := strings.Join(colNames, ",") |
| 35 | tableName := table.Name |
| 36 | sql := fmt.Sprintf(readSQL, cols, pgx.Identifier{tableName}.Sanitize()) |
| 37 | rows, err := c.conn.Query(ctx, sql) |
| 38 | if err != nil { |
| 39 | return err |
| 40 | } |
| 41 | for rows.Next() { |
| 42 | values, err := rows.Values() |
| 43 | if err != nil { |
| 44 | return err |
| 45 | } |
| 46 | |
| 47 | arrowSchema := table.ToArrowSchema() |
| 48 | rb := array.NewRecordBuilder(memory.DefaultAllocator, arrowSchema) |
| 49 | for i := range values { |
| 50 | val, err := prepareValueForResourceSet(arrowSchema.Field(i).Type, values[i]) |
| 51 | if err != nil { |
| 52 | return err |
| 53 | } |
| 54 | s := scalar.NewScalar(arrowSchema.Field(i).Type) |
| 55 | if err := s.Set(val); err != nil { |
| 56 | return err |
| 57 | } |
| 58 | scalar.AppendToBuilder(rb.Field(i), s) |
| 59 | } |
| 60 | res <- rb.NewRecordBatch() |
| 61 | } |
| 62 | rows.Close() |
| 63 | return nil |
| 64 | } |
| 65 | |
| 66 | func prepareValueForResourceSet(dataType arrow.DataType, v any) (any, error) { |
| 67 | if v == nil { |