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

从Google BigTable读取1000万条数据存入Redis的性能瓶颈问题及批量存储优化方案咨询

Fixing High CPU Usage for 10M-Row BigTable-to-Redis Sync

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 03:28:12