如何在Golang中将Spanner行转为JSON/Parquet并定时备份至GCS
解决Spanner快照导出到GCS(Parquet/JSON)的问题
Hey there! Since you're already fetching rows from Spanner, let's build out the rest of the pipeline to convert those rows into your target formats and upload them to GCS. I'll cover both JSON and Parquet, plus the 5-minute scheduling part.
第一步:准备依赖
首先,确保你安装了必要的Go包:
go get cloud.google.com/go/storage go get cloud.google.com/go/spanner go get encoding/json # 对于Parquet,推荐使用Apache Arrow的Parquet库(更贴合行业标准) go get github.com/apache/arrow/go/v14/arrow go get github.com/apache/arrow/go/v14/arrow/array go get github.com/apache/arrow/go/v14/arrow/memory go get github.com/apache/arrow/go/v14/parquet/pqarrow
第二步:将Spanner Rows转换为JSON格式
JSON是相对简单的选项,我们可以把每行数据扫描到一个map[string]interface{},然后收集所有行再序列化为JSON数组:
import ( "context" "encoding/json" "log" "cloud.google.com/go/spanner" "google.golang.org/api/iterator" ) func fetchAndConvertToJSON(ctx context.Context, txn *spanner.ReadWriteTransaction, tableName string, start, end spanner.Timestamp) ([]byte, error) { stmt := spanner.NewStatement("SELECT * FROM " + tableName + " WHERE UpdatedAt >= @startDateTime AND UpdatedAt <= @endDateTime") stmt.Params = map[string]interface{}{ "startDateTime": start, "endDateTime": end, } iter := txn.Query(ctx, stmt) defer iter.Stop() var rows []map[string]interface{} for { row, err := iter.Next() if err == iterator.Done { break } if err != nil { return nil, err } // 动态扫描行数据到map rowMap := make(map[string]interface{}) if err := row.ToStruct(&rowMap); err != nil { log.Printf("Failed to convert row to map: %v", err) continue } rows = append(rows, rowMap) } // 序列化为格式化后的JSON jsonData, err := json.MarshalIndent(rows, "", " ") if err != nil { return nil, err } return jsonData, nil }
第三步:将Spanner Rows转换为Parquet格式
Parquet是列存格式,需要先定义对应的Schema。这里用Apache Arrow的Parquet库来实现,更符合行业标准:
import ( "bytes" "context" "log" "cloud.google.com/go/spanner" "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/pqarrow" "google.golang.org/api/iterator" ) // 定义和你的Spanner表对应的Arrow Schema // 示例:假设表有ID(int64), Name(string), UpdatedAt(timestamp)三个字段 func getTableSchema() *arrow.Schema { return arrow.NewSchema([]arrow.Field{ {Name: "ID", Type: arrow.PrimitiveTypes.Int64, Nullable: false}, {Name: "Name", Type: arrow.BinaryTypes.String, Nullable: true}, {Name: "UpdatedAt", Type: arrow.FixedWidthTypes.Timestamp_ns, Nullable: false}, }, nil) } func fetchAndConvertToParquet(ctx context.Context, txn *spanner.ReadWriteTransaction, tableName string, start, end spanner.Timestamp) ([]byte, error) { stmt := spanner.NewStatement("SELECT * FROM " + tableName + " WHERE UpdatedAt >= @startDateTime AND UpdatedAt <= @endDateTime") stmt.Params = map[string]interface{}{ "startDateTime": start, "endDateTime": end, } iter := txn.Query(ctx, stmt) defer iter.Stop() schema := getTableSchema() mem := memory.NewGoAllocator() builder := array.NewRecordBuilder(mem, schema) defer builder.Release() for { row, err := iter.Next() if err == iterator.Done { break } if err != nil { return nil, err } // 按表结构逐个字段扫描 var id int64 var name *string var updatedAt spanner.Timestamp if err := row.Columns(&id, &name, &updatedAt); err != nil { log.Printf("Failed to scan row: %v", err) continue } // 填充Arrow Record Builder builder.Field(0).(*array.Int64Builder).Append(id) if name != nil { builder.Field(1).(*array.StringBuilder).Append(*name) } else { builder.Field(1).(*array.StringBuilder).AppendNull() } // 将Spanner Timestamp转换为Arrow纳秒级时间戳 builder.Field(2).(*array.TimestampBuilder).Append(updatedAt.Time.UnixNano()) } record := builder.NewRecord() defer record.Release() // 将Record写入Parquet字节流 var buf bytes.Buffer writer, err := pqarrow.NewFileWriter(&buf, schema, pqarrow.DefaultWriterProps()) if err != nil { return nil, err } defer writer.Close() if _, err := writer.Write(record); err != nil { return nil, err } return buf.Bytes(), nil }
注意:Parquet需要严格的Schema定义,你需要根据自己的Spanner表结构调整
getTableSchema和字段扫描逻辑。如果表结构经常变化,可以考虑从Spanner元数据动态生成Schema,但复杂度会更高。
第四步:上传数据到Google Cloud Storage
现在把转换好的字节数据上传到GCS:
import ( "context" "cloud.google.com/go/storage" ) func uploadToGCS(ctx context.Context, bucketName, objectName string, data []byte) error { client, err := storage.NewClient(ctx) if err != nil { return err } defer client.Close() bucket := client.Bucket(bucketName) obj := bucket.Object(objectName) // 创建GCS对象Writer并写入数据 w := obj.NewWriter(ctx) if _, err := w.Write(data); err != nil { return err } // 关闭Writer完成上传 if err := w.Close(); err != nil { return err } return nil }
第五步:实现每5分钟的定时任务
使用time.Ticker定期触发导出任务:
import ( "context" "log" "time" "cloud.google.com/go/spanner" ) func startSnapshotScheduler(ctx context.Context, spannerClient *spanner.Client, tableName, bucketName string) { ticker := time.NewTicker(5 * time.Minute) defer ticker.Stop() // 启动时立即执行一次导出 performSnapshotExport(ctx, spannerClient, tableName, bucketName) for range ticker.C { performSnapshotExport(ctx, spannerClient, tableName, bucketName) } } func performSnapshotExport(ctx context.Context, spannerClient *spanner.Client, tableName, bucketName string) { log.Println("Starting Spanner snapshot export...") // 计算过去5分钟的时间范围 now := time.Now() endTime := spanner.CommitTimestamp(now) startTime := spanner.CommitTimestamp(now.Add(-5 * time.Minute)) // 选择JSON或Parquet格式,这里示例用JSON data, err := fetchAndConvertToJSON(ctx, spannerClient.ReadWriteTransaction(ctx), tableName, startTime, endTime) // 如果用Parquet,替换为下面一行 // data, err := fetchAndConvertToParquet(ctx, spannerClient.ReadWriteTransaction(ctx), tableName, startTime, endTime) if err != nil { log.Printf("Failed to fetch and convert data: %v", err) return } // 生成带时间戳的唯一文件名 objectName := "spanner-snapshots/" + tableName + "_" + now.Format("20060102_150405") + ".json" // Parquet格式的文件名: // objectName := "spanner-snapshots/" + tableName + "_" + now.Format("20060102_150405") + ".parquet" if err := uploadToGCS(ctx, bucketName, objectName, data); err != nil { log.Printf("Failed to upload to GCS: %v", err) return } log.Printf("Successfully exported snapshot to gs://%s/%s", bucketName, objectName) }
最后:主函数整合
把所有部分整合到主函数中:
import ( "context" "log" "cloud.google.com/go/spanner" ) func main() { ctx := context.Background() // 替换为你的Spanner实例和数据库路径 spannerDB := "projects/your-project-id/instances/your-instance-id/databases/your-db-id" // 替换为你的GCS Bucket名称 bucketName := "your-gcs-bucket-name" tableName := "your-table-name" spannerClient, err := spanner.NewClient(ctx, spannerDB) if err != nil { log.Fatalf("Failed to create Spanner client: %v", err) } defer spannerClient.Close() startSnapshotScheduler(ctx, spannerClient, tableName, bucketName) }
额外建议:
- 错误处理:可以添加重试逻辑,比如GCS上传失败时重试2-3次
- 性能优化:如果数据量很大,不要一次性加载所有行到内存,可以考虑流式写入GCS(直接将Parquet/JSON写入流和GCS Writer对接)
- 权限:确保你的服务账号拥有Spanner读取权限和GCS写入权限
内容的提问来源于stack exchange,提问作者Bruce Mu
相关产品推荐
相关产品推荐

