googleapis / googleapis/google-cloud-go

bigquery:storage write client doesn't work for arrow format

Open
#12,478 4 comments 0 reactions 1 assignee Claimed by @alvarowolfx View on GitHub
api: bigquery priority: p3
Dominant language
Go
Stars
4.5k
Forks
1.6k
Avg merge
1d 13h
Merged PRs (30d)
109

Description

## Client

## Environment

go version go1.23.4 darwin/arm64
Tried with arrow12 to arrow18

## Code and Dependencies

* Following is the code I ran which isn't working

```go
import (
arrow15 "github.com/apache/arrow/go/v15/arrow"
arrow15array "github.com/apache/arrow/go/v15/arrow/array"
arrow15ipc "github.com/apache/arrow/go/v15/arrow/ipc"
arrow15memory "github.com/apache/arrow/go/v15/arrow/memory"
)

func writeToBQClient(record arrow15.Record, writeStreamName string, appendStream storagepb.BigQueryWrite_AppendRowsClient) error {

mem := arrow15memory.NewGoAllocator()
defer mem.Close()

// Serialize schema using IPC format (just the schema, no data)
var schemaBuf bytes.Buffer
schemaWriter := arrow15ipc.NewWriter(&schemaBuf, arrow15ipc.WithSchema(record.Schema()), arrow15ipc.WithAllocator(mem))
err := schemaWriter.Close() // This writes just the schema without any records
if err != nil {
log.Logger().WithError(err).Error("Failed to serialize schema")
return err
}
schemaData := schemaBuf.Bytes()

// Serialize record batch using IPC format
var recordBuf bytes.Buffer
recordWriter := arrow15ipc.NewWriter(&recordBuf, arrow15ipc.WithSchema(record.Schema()), arrow15ipc.WithAllocator(mem))

err = recordWriter.Write(record)
if err != nil {
log.Logger().WithError(err).Error("Failed to write Arrow record")
recordWriter.Close()
return err
}

err = recordWriter.Close()
if err != nil {
log.Logger().WithError(err).Error("Failed to close Arrow writer")
return err
}

recordData := recordBuf.Bytes()

log.Logger().Infof("Schema IPC data size: %d bytes", len(schemaData))
log.Logger().Infof("Record IPC data size: %d bytes", len(recordData))

request := &storagepb.AppendRowsRequest{
WriteStream: writeStreamName,
Rows: &storagepb.AppendRowsRequest_ArrowRows{
ArrowRows: &storagepb.AppendRowsRequest_ArrowData{
WriterSchema: &storagepb.ArrowSchema{
SerializedSchema: schemaData,
},
Rows: &storagepb.ArrowRecordBatch{
SerializedRecordBatch: recordData,
},
},
},
}

// Send the request
err = appendStream.Send(request)
if err != nil {
log.Logger().WithError(err).Error("Failed to send AppendRows request")
return err
}

err = appendStream.CloseSend()
if err != nil {
log.Logger().WithError(err).Error("Failed to close AppendRows request")
return err
}

return nil
}

func Main() {
// Create Arrow v12 record with single column c0 of type string and value "tushartg"
pool := arrow15memory.NewGoAllocator()

// Create schema with single string column
schema := arrow15.NewSchema([]arrow15.Field{
{Name: "c0", Type: arrow15.BinaryTypes.String, Nullable: true},
}, nil)

// Create record builder
builder := arrow15array.NewRecordBuilder(pool, schema)
defer builder.Release()

// Add single row with value "tushartg"
builder.Field(0).(*arrow15array.StringBuilder).Append("tushartg")

// Build the record
record := builder.NewRecord()
defer record.Release()

storageWriteClient, err := bqstorage.NewBigQueryWriteClient(ctx, c.authOption)
if err != nil {
log.Logger().WithError(err).Error("Failed to create fresh BigQuery Storage write client")
return err
}
defer storageWriteClient.Close()

stream, err := storageWriteClient.CreateWriteStream(ctx, &storagepb.CreateWriteStreamRequest{
Parent: fmt.Sprintf("projects/%s/datasets/%s/tables/%s", c.projectID, schema, tableName),
WriteStream: &storagepb.WriteStream{
Type: storagepb.WriteStream_COMMITTED,
},
})
if err != nil {
return err
}

appendStream, err := storageWriteClient.AppendRows(ctx)
if err != nil {
log.Logger().WithError(err).Error("Failed to create AppendRows stream")
return err
}

// Write to BigQuery and handle response synchronously
err = writeToBQClient(record, stream.Name, appendStream)
if err != nil {
log.Logger().WithError(err).Error("Failed to write to BigQuery client")
return err
}

// Receive responses
for {
resp, err := appendStream.Recv()
if err == io.EOF {
break
}
if err != nil {
log.Logger().WithError(err).Error("Failed to receive AppendRows response")
return err
}
if resp.GetError() != nil {
log.Logger().Errorf("BigQuery error: %v", resp.GetError())
return errors.NewErrorf(errors.Internal, "BigQuery error: %v", resp.GetError())
}
}
}
```

Error received
```
ERROR: BigQuery error: code:3 message:"Header-type of flatbuffer-encoded Message is not RecordBatch. Entity: projects/PROJECT/datasets/DATASET/tables/TABLENAME/streams/_default" details:{[type.googleapis.com/google.rpc.DebugInfo]:{detail:"INVALID_ARGUMENT: Header-type of flatbuffer-encoded Message is not RecordBatch. [type.googleapis.com/util.MessageSetPayload='[cloud.helix.vortex.VortexStatus] { error_code: CONVERTER_PARSE_ERROR component: CONVERTER is_public_error: true }']"}}
```

## Expected behavior

Data gets ingested into table

## Actual behavior

Error: `Header-type of flatbuffer-encoded Message is not RecordBatch`

## Additional context

I want to know how to serialize the arrow records basically. Can you please help improve documentation?

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.