Golang gRPC服务端动态查询结果转Parquet并发送至客户端实现咨询
Go实现gRPC服务端动态查询转Parquet并发送给客户端
一、定义gRPC接口
针对动态生成的Parquet数据,推荐用流式接口避免单次传输过大,先编写.proto文件:
syntax = "proto3"; package queryservice; service QueryParquetService { // 发送查询请求,返回流式Parquet分片数据 rpc ExecuteQuery(QueryRequest) returns (stream ParquetChunk); } message QueryRequest { string sql = 1; // 客户端传入的动态查询SQL } message ParquetChunk { bytes data = 1; // 分片的Parquet二进制数据 }
用protoc生成Go代码:
protoc --go_out=. --go-grpc_out=. queryservice.proto
二、服务端核心实现
1. 执行动态查询并提取元数据与数据
以MySQL为例,用database/sql实现动态查询,同时获取列名和数据行:
import ( "database/sql" _ "github.com/go-sql-driver/mysql" ) // 执行SQL,返回列名列表和每行数据的map集合 func executeDynamicQuery(db *sql.DB, sqlStr string) ([]string, []map[string]interface{}, error) { rows, err := db.Query(sqlStr) if err != nil { return nil, nil, err } defer rows.Close() // 获取查询结果的列名 cols, err := rows.Columns() if err != nil { return nil, nil, err } // 遍历所有数据行,转为map存储 var result []map[string]interface{} for rows.Next() { values := make([]interface{}, len(cols)) valuePtrs := make([]interface{}, len(cols)) for i := range cols { valuePtrs[i] = &values[i] } if err := rows.Scan(valuePtrs...); err != nil { return nil, nil, err } rowMap := make(map[string]interface{}) for i, col := range cols { rowMap[col] = values[i] } result = append(result, rowMap) } return cols, result, rows.Err() }
2. 动态生成Parquet Schema并转换数据
使用Apache Arrow的Go库处理Parquet,它对动态Schema支持更完善:
go get github.com/apache/arrow/go/v14/arrow go get github.com/apache/arrow/go/v14/parquet/pqarrow
实现转换逻辑:
import ( "bytes" "fmt" "github.com/apache/arrow/go/v14/arrow" "github.com/apache/arrow/go/v14/arrow/array" "github.com/apache/arrow/go/v14/arrow/memory" "github.com/apache/arrow/go/v14/parquet/file" "github.com/apache/arrow/go/v14/parquet/pqarrow" ) // 将动态查询结果转为Parquet字节流 func convertToParquet(cols []string, rows []map[string]interface{}) ([]byte, error) { // 动态生成Arrow Schema(Parquet依赖Arrow Schema定义结构) fields := make([]arrow.Field, len(cols)) for i, col := range cols { // 这里简化类型映射,实际需根据数据库列类型对应Arrow类型(如int64、float64等) fields[i] = arrow.Field{Name: col, Type: arrow.BinaryTypes.String, Nullable: true} } schema := arrow.NewSchema(fields, nil) // 创建Parquet写入器 mem := memory.NewAllocator() defer mem.Release() buf := bytes.NewBuffer(nil) writer, err := pqarrow.NewWriter(buf, schema, pqarrow.WithWriterProps(file.WithCompression(file.Snappy))) if err != nil { return nil, err } defer writer.Close() // 构造Arrow记录 builders := array.NewBuilders(mem, schema) defer builders.Release() for _, row := range rows { for i, col := range cols { val := row[col] // 根据实际类型写入builder,示例统一转为字符串,实际需适配不同类型 switch v := val.(type) { case string: builders[i].(*array.StringBuilder).Append(v) case []byte: builders[i].(*array.StringBuilder).Append(string(v)) case nil: builders[i].AppendNull() default: builders[i].(*array.StringBuilder).Append(fmt.Sprintf("%v", v)) } } } record := array.NewRecord(schema, builders.NewArrays(), int64(len(rows))) defer record.Release() // 写入Parquet数据 if _, err := writer.Write(record); err != nil { return nil, err } return buf.Bytes(), nil }
3. gRPC服务端整合逻辑
将查询、转换、流式发送整合到gRPC服务中:
import ( "context" "net" "google.golang.org/grpc" ) type queryParquetServer struct { UnimplementedQueryParquetServiceServer db *sql.DB } func (s *queryParquetServer) ExecuteQuery(req *QueryRequest, stream QueryParquetService_ExecuteQueryServer) error { // 执行动态查询 cols, rows, err := executeDynamicQuery(s.db, req.Sql) if err != nil { return err } // 转换为Parquet字节流 parquetData, err := convertToParquet(cols, rows) if err != nil { return err } // 分片发送(1MB每片,避免单次传输过大) chunkSize := 1024 * 1024 for i := 0; i < len(parquetData); i += chunkSize { end := i + chunkSize if end > len(parquetData) { end = len(parquetData) } if err := stream.Send(&ParquetChunk{Data: parquetData[i:end]}); err != nil { return err } } return nil } // 启动gRPC服务 func startGRPCServer(db *sql.DB, addr string) error { s := grpc.NewServer() RegisterQueryParquetServiceServer(s, &queryParquetServer{db: db}) lis, err := net.Listen("tcp", addr) if err != nil { return err } return s.Serve(lis) }
三、客户端接收与解析
客户端接收流式分片,合并后解析Parquet数据:
import ( "bytes" "context" "io" "github.com/apache/arrow/go/v14/arrow/memory" "github.com/apache/arrow/go/v14/parquet/file" "github.com/apache/arrow/go/v14/parquet/pqarrow" "google.golang.org/grpc" ) func fetchParquetData(ctx context.Context, conn *grpc.ClientConn, sql string) ([]map[string]interface{}, error) { client := NewQueryParquetServiceClient(conn) stream, err := client.ExecuteQuery(ctx, &QueryRequest{Sql: sql}) if err != nil { return nil, err } // 合并所有Parquet分片 var parquetBuf bytes.Buffer for { chunk, err := stream.Recv() if err == io.EOF { break } if err != nil { return nil, err } parquetBuf.Write(chunk.Data) } // 解析Parquet数据 reader, err := file.NewReader(bytes.NewReader(parquetBuf.Bytes()), file.WithReadProps(nil)) if err != nil { return nil, err } defer reader.Close() arrowReader, err := pqarrow.NewFileReader(reader, pqarrow.WithAllocator(memory.NewAllocator())) if err != nil { return nil, err } defer arrowReader.Release() // 读取所有记录并转为map集合 var result []map[string]interface{} for arrowReader.Next() { record := arrowReader.Record() defer record.Release() cols := record.Schema().Fields() for rowIdx := 0; rowIdx < int(record.NumRows()); rowIdx++ { rowMap := make(map[string]interface{}) for colIdx, col := range cols { rowMap[col.Name] = record.Column(colIdx).Value(rowIdx) } result = append(result, rowMap) } } return result, arrowReader.Err() }
关键注意事项
- 类型映射:实际项目中需根据数据库列类型(int、float、datetime等)准确映射到Arrow对应的类型,避免数据失真。
- 内存优化:若查询结果量极大,建议服务端边查询边写入Parquet并流式发送,避免内存溢出。
- 错误处理:示例简化了错误处理,实际需完善异常捕获和兜底逻辑。
内容的提问来源于stack exchange,提问作者sama
相关产品推荐
相关产品推荐

