googleapis / googleapis/google-cloud-go
bigquery:storage write client doesn't work for arrow format
- 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
Assessment
This issue has not been assessed yet.