refactor(aws-sdk)!: migrate to transfermanager

Signed-off-by: MadDogOwner <xiaoran@xrgzs.top>
This commit is contained in:
MadDogOwner
2026-07-16 22:08:57 +08:00
parent f99a1f9a50
commit 80487122ea
9 changed files with 51 additions and 58 deletions
+4 -8
View File
@@ -21,7 +21,7 @@ import (
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
"github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager"
"github.com/go-resty/resty/v2"
log "github.com/sirupsen/logrus"
)
@@ -219,19 +219,15 @@ func (d *Pan123) Put(ctx context.Context, dstDir model.Obj, file model.FileStrea
o.BaseEndpoint = aws.String(resp.Data.EndPoint)
o.UsePathStyle = true
})
uploader := manager.NewUploader(s3Client)
if file.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize {
uploader.PartSize = file.GetSize() / int64(manager.MaxUploadParts-1)
}
input := &s3.PutObjectInput{
tmClient := transfermanager.New(s3Client)
_, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{
Bucket: &resp.Data.Bucket,
Key: &resp.Data.Key,
Body: driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{
Reader: file,
UpdateProgress: up,
}),
}
_, err = uploader.Upload(ctx, input)
})
if err != nil {
return err
}
+9 -7
View File
@@ -15,11 +15,12 @@ import (
"github.com/OpenListTeam/OpenList/v4/internal/model"
"github.com/OpenListTeam/OpenList/v4/internal/op"
"github.com/OpenListTeam/OpenList/v4/internal/stream"
"github.com/OpenListTeam/OpenList/v4/pkg/utils"
"github.com/OpenListTeam/OpenList/v4/pkg/http_range"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
"github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager"
"github.com/city404/v6-public-rpc-proto/go/v6/common"
pbPublicUser "github.com/city404/v6-public-rpc-proto/go/v6/user"
pubUserFile "github.com/city404/v6-public-rpc-proto/go/v6/userfile"
@@ -381,14 +382,15 @@ func (d *HalalCloud) put(ctx context.Context, dstDir model.Obj, fileStream model
o.BaseEndpoint = aws.String(result.Endpoint)
o.UsePathStyle = true
})
uploader := manager.NewUploader(s3Client, func(u *manager.Uploader) {
u.Concurrency = d.uploadThread
tmClient := transfermanager.New(s3Client, func(o *transfermanager.Options) {
if fileStream.GetSize() > int64(8*utils.MB)*10000 {
o.PartSizeBytes = fileStream.GetSize() / 9999
}
o.Concurrency = d.uploadThread
})
if fileStream.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize {
uploader.PartSize = fileStream.GetSize() / int64(manager.MaxUploadParts-1)
}
reader := driver.NewLimitedUploadStream(ctx, fileStream)
_, err = uploader.Upload(ctx, &s3.PutObjectInput{
_, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{
Bucket: aws.String(result.Bucket),
Key: aws.String(result.Key),
Body: io.TeeReader(reader, driver.NewProgress(fileStream.GetSize(), up)),
+9 -8
View File
@@ -17,7 +17,7 @@ import (
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
"github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager"
"github.com/go-resty/resty/v2"
"github.com/google/uuid"
log "github.com/sirupsen/logrus"
@@ -184,11 +184,13 @@ func (d *MediaTrack) Put(ctx context.Context, dstDir model.Obj, file model.FileS
if err != nil {
return err
}
uploader := manager.NewUploader(s3Client)
if file.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize {
uploader.PartSize = file.GetSize() / int64(manager.MaxUploadParts-1)
}
input := &s3.PutObjectInput{
tmClient := transfermanager.New(s3Client, func(o *transfermanager.Options) {
if file.GetSize() > int64(8*utils.MB)*10000 {
o.PartSizeBytes = file.GetSize() / 9999
}
})
_, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{
Bucket: &resp.Data.Bucket,
Key: &resp.Data.Object,
Body: driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{
@@ -198,8 +200,7 @@ func (d *MediaTrack) Put(ctx context.Context, dstDir model.Obj, file model.FileS
},
UpdateProgress: up,
}),
}
_, err = uploader.Upload(ctx, input)
})
if err != nil {
return err
}
+8 -8
View File
@@ -20,7 +20,7 @@ import (
"github.com/aws/smithy-go"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/service/s3/types"
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
"github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager"
"github.com/pkg/errors"
log "github.com/sirupsen/logrus"
)
@@ -212,14 +212,15 @@ func (d *S3) Remove(ctx context.Context, obj model.Obj) error {
}
func (d *S3) Put(ctx context.Context, dstDir model.Obj, s model.FileStreamer, up driver.UpdateProgress) error {
uploader := manager.NewUploader(d.client)
if s.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize {
uploader.PartSize = s.GetSize() / int64(manager.MaxUploadParts-1)
}
key := getKey(stdpath.Join(dstDir.GetPath(), s.GetName()), false)
contentType := s.GetMimetype()
log.Debugln("key:", key)
input := &s3.PutObjectInput{
tmClient := transfermanager.New(d.client, func(o *transfermanager.Options) {
if s.GetSize() > int64(8*utils.MB)*10000 {
o.PartSizeBytes = s.GetSize() / 9999
}
})
_, err := tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{
Bucket: &d.Bucket,
Key: &key,
Body: driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{
@@ -227,8 +228,7 @@ func (d *S3) Put(ctx context.Context, dstDir model.Obj, s model.FileStreamer, up
UpdateProgress: up,
}),
ContentType: &contentType,
}
_, err := uploader.Upload(ctx, input)
})
return err
}
+9 -8
View File
@@ -17,7 +17,7 @@ import (
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
"github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager"
"github.com/go-resty/resty/v2"
log "github.com/sirupsen/logrus"
)
@@ -248,11 +248,13 @@ func (d *Teambition) newUpload(ctx context.Context, dstDir model.Obj, stream mod
if err != nil {
return err
}
uploader := manager.NewUploader(s3Client)
if stream.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize {
uploader.PartSize = stream.GetSize() / int64(manager.MaxUploadParts-1)
}
input := &s3.PutObjectInput{
tmClient := transfermanager.New(s3Client, func(o *transfermanager.Options) {
if stream.GetSize() > int64(8*utils.MB)*10000 {
o.PartSizeBytes = stream.GetSize() / 9999
}
})
_, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{
Bucket: &uploadToken.Upload.Bucket,
Key: &uploadToken.Upload.Key,
ContentDisposition: &uploadToken.Upload.ContentDisposition,
@@ -261,8 +263,7 @@ func (d *Teambition) newUpload(ctx context.Context, dstDir model.Obj, stream mod
Reader: stream,
UpdateProgress: up,
}),
}
_, err = uploader.Upload(ctx, input)
})
if err != nil {
return err
}
+4 -6
View File
@@ -19,7 +19,7 @@ import (
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
"github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager"
"github.com/go-resty/resty/v2"
)
@@ -416,11 +416,9 @@ func (xc *XunLeiCommon) Put(ctx context.Context, dstDir model.Obj, file model.Fi
}, func(o *s3.Options) {
o.BaseEndpoint = aws.String(param.Endpoint)
})
uploader := manager.NewUploader(s3Client)
if file.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize {
uploader.PartSize = file.GetSize() / int64(manager.MaxUploadParts-1)
}
_, err = uploader.Upload(ctx, &s3.PutObjectInput{
tmClient := transfermanager.New(s3Client)
_, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{
Bucket: aws.String(param.Bucket),
Key: aws.String(param.Key),
Body: driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{
+4 -6
View File
@@ -21,7 +21,7 @@ import (
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
"github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager"
"github.com/go-resty/resty/v2"
)
@@ -530,11 +530,9 @@ func (xc *XunLeiBrowserCommon) Put(ctx context.Context, dstDir model.Obj, stream
}, func(o *s3.Options) {
o.BaseEndpoint = aws.String(param.Endpoint)
})
uploader := manager.NewUploader(s3Client)
if stream.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize {
uploader.PartSize = stream.GetSize() / int64(manager.MaxUploadParts-1)
}
_, err = uploader.Upload(ctx, &s3.PutObjectInput{
tmClient := transfermanager.New(s3Client)
_, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{
Bucket: aws.String(param.Bucket),
Key: aws.String(param.Key),
Expires: aws.Time(param.Expiration),
+4 -6
View File
@@ -20,7 +20,7 @@ import (
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/credentials"
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/feature/s3/manager"
"github.com/aws/aws-sdk-go-v2/feature/s3/transfermanager"
"github.com/go-resty/resty/v2"
)
@@ -406,11 +406,9 @@ func (xc *XunLeiXCommon) Put(ctx context.Context, dstDir model.Obj, file model.F
}, func(o *s3.Options) {
o.BaseEndpoint = aws.String(param.Endpoint)
})
uploader := manager.NewUploader(s3Client)
if file.GetSize() > int64(manager.MaxUploadParts)*manager.DefaultUploadPartSize {
uploader.PartSize = file.GetSize() / int64(manager.MaxUploadParts-1)
}
_, err = uploader.Upload(ctx, &s3.PutObjectInput{
tmClient := transfermanager.New(s3Client)
_, err = tmClient.UploadObject(ctx, &transfermanager.UploadObjectInput{
Bucket: aws.String(param.Bucket),
Key: aws.String(param.Key),
Body: driver.NewLimitedUploadStream(ctx, &driver.ReaderUpdatingProgress{
-1
View File
@@ -200,7 +200,6 @@ require (
github.com/andreburgaud/crypt2go v1.8.0 // indirect
github.com/andybalholm/brotli v1.2.2 // indirect
github.com/aws/aws-sdk-go-v2/credentials v1.19.29
github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.22.13
github.com/aws/aws-sdk-go-v2/service/s3 v1.105.1
github.com/aws/smithy-go v1.27.3
github.com/axgle/mahonia v0.0.0-20180208002826-3358181d7394