从Google BigTable读取1000万条数据存入Redis的性能瓶颈问题及批量存储优化方案咨询
Let’s break down exactly why your code is spiking CPU with 10 million records, and walk through targeted fixes—starting with the Redis bottleneck, which is the biggest offender here.
1. Ditch Per-Record Goroutines for Redis Batch Pipelines
Your current approach spins up a goroutine for every single Redis HMSET/EXPIRE call—10 million goroutines total. This creates massive context-switching overhead, which is why your CPU is maxing out. Instead, use Redis pipelines to send batches of commands in one network round-trip.
Here’s a revised saveToRedis function with batching and pipelines:
func saveToRedis(redisClient *redisv2.Connectionx, rowData []class_ad.RedisData, now time.Time) { const batchSize = 1000 // Tune this based on your Redis server's capacity (start with 500-2000) totalRecords := len(rowData) success, fail := 0, 0 for i := 0; i < totalRecords; i += batchSize { endIdx := min(i+batchSize, totalRecords) batch := rowData[i:endIdx] // Initialize a pipeline to queue commands pipe := redisClient.Pipeline() for _, data := range batch { // Queue HMSET command hmsetErr := pipe.HMSet(data.Key, map[string]string{ "COLUMN_1": strconv.FormatFloat(data.Column1, 'f', -1, 64), "COLUMN_2": strconv.FormatFloat(data.Conversion, 'f', -1, 64), "CREATION_DATE": now.String(), "MODIFICATION_DATE": now.String(), }) if hmsetErr != nil { fail++ continue } // Queue EXPIREAT command alongside HMSET expireErr := pipe.ExpireAt(data.Key, endDate.Unix()) if expireErr != nil { fail++ continue } success++ } // Execute all queued commands in one go _, execErr := pipe.Exec(context.Background()) if execErr != nil { log.Printf("Pipeline failed for batch %d-%d: %v", i, endIdx, execErr) // Adjust failure count for the entire failed batch fail += len(batch) - (success - (len(batch) - fail)) } } log.Printf("Sync complete: Success %v, Failure %v", success, fail) }
This cuts network round-trips from 10 million to ~10,000, eliminating the goroutine context-switching nightmare.
2. Preallocate Slice Capacity to Avoid Memory Copy Overhead
Your redisRd slice starts empty and gets appended to millions of times. Each append can trigger a memory resize (doubling capacity each time), which involves copying the entire slice to new memory—this is a huge CPU drain. Preallocate the slice with an estimated capacity upfront:
// Preallocate capacity for 10M records (adjust if you know a more precise number) redisRd := make([]class_ad.RedisData, 0, 10_000_000) err = tbl.ReadRows( context.Background(), bt.PrefixRange(currDay), func(row bt.Row) bool { data, err := sS.readAllRowData(row) // Stop ignoring errors! if err != nil { log.Printf("Failed to read row %s: %v", row.Key(), err) return true } rd, err := sS.parseAllRows(data) if err != nil { log.Printf("Failed to parse row %s: %v", row.Key(), err) return true } redisRd = append(redisRd, rd...) return true }, )
Preallocating eliminates expensive memory copies, which will noticeably reduce CPU usage for large datasets.
3. Use a Redis Connection Pool (Not a Single Connection)
Using one Redis connection for all operations creates a bottleneck—even if the client is thread-safe, a single connection can’t handle 10 million commands efficiently. Switch to a connection pool to distribute load across multiple connections:
// Initialize the pool once (e.g., in your struct's constructor) redisPool := redisv2.NewPool(&redisv2.Options{ Addr: "your-redis-host:6379", MaxIdle: 10, // Keep idle connections warm MaxActive: 50, // Adjust based on your Redis server's capacity }) // In saveToRedis, get a connection from the pool conn, err := redisPool.Get() if err != nil { log.Printf("Failed to get Redis connection: %v", err) return } defer conn.Close() // Use this connection for your pipeline batches
A pool lets you parallelize batch operations, reducing wait time and CPU load from blocked I/O.
4. Process BigTable Rows in Batches (Don’t Load All into Memory)
Holding 10 million records in memory at once is unnecessary and adds memory management overhead. Instead, process batches of rows as you read them from BigTable, overlapping read and write operations:
const batchSize = 1000 var tempBatch []class_ad.RedisData = make([]class_ad.RedisData, 0, batchSize) err = tbl.ReadRows( context.Background(), bt.PrefixRange(currDay), func(row bt.Row) bool { data, err := sS.readAllRowData(row) if err != nil { log.Printf("Row read error: %v", err) return true } rd, err := sS.parseAllRows(data) if err != nil { log.Printf("Row parse error: %v", err) return true } tempBatch = append(tempBatch, rd...) // When tempBatch hits size, send to Redis immediately if len(tempBatch) >= batchSize { saveToRedis(redisClient, tempBatch, now) tempBatch = tempBatch[:0] // Reset slice (keeps preallocated capacity) } return true }, ) // Save any remaining rows in the final batch if len(tempBatch) > 0 { saveToRedis(redisClient, tempBatch, now) }
This reduces memory pressure and lets you start writing to Redis before all BigTable data is read, improving overall throughput.
5. Precompute String Conversions (Don’t Do It During Redis Writes)
Your current code converts floats to strings during Redis operations, which adds unnecessary CPU load during the time-sensitive write phase. Move these conversions to the parseAllRows step:
// Update your RedisData struct to store precomputed strings type RedisData struct { Key string Column1Str string ConversionStr string // ... other fields } func (sS *someStruct) parseAllRows(data []ColumnData) ([]class_ad.RedisData, error) { var rd class_ad.RedisData for _, v := range data { if v.ColumnName == "SomeColumn1" { col1Val, err := strconv.ParseFloat(string(v.Value), 64) if err != nil { return nil, fmt.Errorf("parse SomeColumn1: %v", err) } rd.Column1Str = strconv.FormatFloat(col1Val, 'f', -1, 64) continue } if v.ColumnName == "SomeColumn2" { convVal, err := strconv.ParseFloat(string(v.Value), 64) if err != nil { return nil, fmt.Errorf("parse SomeColumn2: %v", err) } rd.ConversionStr = strconv.FormatFloat(convVal, 'f', -1, 64) continue } // ... handle other columns } return []class_ad.RedisData{rd}, nil } // Then in saveToRedis, use the precomputed strings hmsetErr := pipe.HMSet(data.Key, map[string]string{ "COLUMN_1": data.Column1Str, "COLUMN_2": data.ConversionStr, // ... })
This shifts CPU work to the parsing phase, leaving Redis operations as lightweight as possible.
Start with the Redis pipeline and batching changes first—those will give you the biggest immediate performance boost. Combine them with the other optimizations, and you should see a dramatic drop in CPU usage while handling your 10 million daily records.
内容的提问来源于stack exchange,提问作者Kunal Vohra

