You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 08:47:46