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

Golang调用DynamoDB PutItem插入缓慢问题求助

优化Golang向DynamoDB批量写入的方案

兄弟,你猜的完全没错——默认情况下你调用的PutItem是同步阻塞的,每一条数据都要等和DynamoDB的网络请求往返完成才会处理下一条,这才是速度暴跌的核心原因。毕竟本地解析JSON是CPU操作,而DynamoDB API调用是跨网络的IO操作,耗时根本不在一个量级上。我之前做过Golang向DynamoDB批量插数据的优化,给你分享几个实用方案:

1. 优先使用BatchWriteItem批量写入API

DynamoDB提供的BatchWriteItem允许你一次向多个表写入最多25条数据(每个表的条目数不能超过25),这能大幅减少网络请求的次数,直接提升吞吐量。需要注意的是,这个API可能会返回UnprocessedItems(比如因为瞬时负载过高或者其他原因没处理成功的条目),所以你需要循环重试这些未处理的项,直到全部完成。

简单的代码示例:

import (
    "github.com/aws/aws-sdk-go/aws"
    "github.com/aws/aws-sdk-go/aws/session"
    "github.com/aws/aws-sdk-go/service/dynamodb"
    "github.com/aws/aws-sdk-go/service/dynamodb/dynamodbattribute"
)

func batchWriteItems(svc *dynamodb.DynamoDB, tableName string, items []interface{}) error {
    var requests []*dynamodb.WriteRequest
    for _, item := range items {
        av, err := dynamodbattribute.MarshalMap(item)
        if err != nil {
            return err
        }
        requests = append(requests, &dynamodb.WriteRequest{
            PutRequest: &dynamodb.PutRequest{Item: av},
        })
    }

    // 拆分批次,每个批次最多25条
    for i := 0; i < len(requests); i += 25 {
        end := i + 25
        if end > len(requests) {
            end = len(requests)
        }
        batch := requests[i:end]

        input := &dynamodb.BatchWriteItemInput{
            RequestItems: map[string][]*dynamodb.WriteRequest{
                tableName: batch,
            },
        }

        resp, err := svc.BatchWriteItem(input)
        if err != nil {
            return err
        }

        // 处理未完成的条目,循环重试
        for len(resp.UnprocessedItems) > 0 {
            input = &dynamodb.BatchWriteItemInput{
                RequestItems: resp.UnprocessedItems,
            }
            resp, err = svc.BatchWriteItem(input)
            if err != nil {
                return err
            }
        }
    }
    return nil
}

2. 用Goroutine并发执行写入请求

如果你的业务场景不适合用批量API(比如每条数据的处理逻辑不一样),可以用Goroutine来并发执行PutItem请求,但一定要控制并发数,避免一下子发起太多请求导致不必要的资源消耗(哪怕WCU足够,过多的并发也可能带来网络层面的开销)。

比如用带缓冲的Channel做并发控制:

import (
    "fmt"
    "github.com/aws/aws-sdk-go/aws"
    "github.com/aws/aws-sdk-go/aws/session"
    "github.com/aws/aws-sdk-go/service/dynamodb"
    "github.com/aws/aws-sdk-go/service/dynamodb/dynamodbattribute"
)

func concurrentPutItems(svc *dynamodb.DynamoDB, tableName string, items []interface{}, concurrency int) error {
    errChan := make(chan error, len(items))
    itemChan := make(chan interface{}, len(items))

    // 把所有数据放到channel里
    for _, item := range items {
        itemChan <- item
    }
    close(itemChan)

    // 启动指定数量的goroutine
    for i := 0; i < concurrency; i++ {
        go func() {
            for item := range itemChan {
                av, err := dynamodbattribute.MarshalMap(item)
                if err != nil {
                    errChan <- err
                    return
                }
                input := &dynamodb.PutItemInput{
                    TableName: aws.String(tableName),
                    Item:      av,
                }
                _, err = svc.PutItem(input)
                if err != nil {
                    errChan <- err
                    return
                }
            }
            errChan <- nil
        }()
    }

    // 收集错误
    var errs []error
    for i := 0; i < concurrency; i++ {
        if err := <-errChan; err != nil {
            errs = append(errs, err)
        }
    }
    close(errChan)

    if len(errs) > 0 {
        return fmt.Errorf("some put operations failed: %v", errs)
    }
    return nil
}

建议根据你的EC2实例配置和DynamoDB的WCU设置,调整concurrency参数(比如从100开始测试,逐步找到最优值)。

3. 优化DynamoDB客户端的HTTP配置

默认的AWS SDK HTTP客户端配置可能不够高效,你可以调整连接池的大小,减少TCP连接建立的开销:

import (
    "net/http"
    "time"
    "github.com/aws/aws-sdk-go/aws"
    "github.com/aws/aws-sdk-go/aws/session"
    "github.com/aws/aws-sdk-go/service/dynamodb"
)

func createOptimizedDynamoDBClient() *dynamodb.DynamoDB {
    httpClient := &http.Client{
        Transport: &http.Transport{
            MaxIdleConns:        200, // 增大空闲连接数
            MaxIdleConnsPerHost: 100, // 每个主机的最大空闲连接数
            IdleConnTimeout:     60 * time.Second, // 空闲连接超时时间
            TLSHandshakeTimeout: 10 * time.Second,
        },
        Timeout: 30 * time.Second, // 请求超时时间
    }

    sess := session.Must(session.NewSession(&aws.Config{
        Region:     aws.String("your-region"),
        HTTPClient: httpClient,
    }))
    return dynamodb.New(sess)
}

4. 事务写入(如果需要原子性)

如果你的业务要求写入两个DynamoDB表的操作必须原子性(要么都成功,要么都失败),可以用TransactWriteItems API,它允许你在一次请求中完成多个表的写入操作,比单独调用两次PutItem高效得多,因为只需要一次网络往返。

总结

优先尝试BatchWriteItem,它能最大程度减少网络请求次数;如果需要更灵活的并发控制,结合Goroutine和连接池优化;如果有原子性需求,用事务写入。按照这些方案优化后,你的写入速度应该能提升几十倍甚至上百倍,接近你本地解析的速度水平。

内容的提问来源于stack exchange,提问作者WeCanBeFriends

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:49:46