MCPcopy Create free account
hub / github.com/cloudquery/cloudquery / Read

Method Read

plugins/destination/postgresql/client/read.go:25–64  ·  view source on GitHub ↗
(ctx context.Context, table *schema.Table, res chan<- arrow.RecordBatch)

Source from the content-addressed store, hash-verified

23)
24
25func (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
66func prepareValueForResourceSet(dataType arrow.DataType, v any) (any, error) {
67 if v == nil {

Callers 1

Calls 5

NextMethod · 0.80
QueryMethod · 0.65
CloseMethod · 0.65
SetMethod · 0.45

Tested by 1