Golang调用DynamoDB PutItem插入缓慢问题求助
兄弟,你猜的完全没错——默认情况下你调用的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

