diff --git a/cmd/digger/subcmd/s3.go b/cmd/digger/subcmd/s3.go index 7a803f1..e339e96 100644 --- a/cmd/digger/subcmd/s3.go +++ b/cmd/digger/subcmd/s3.go @@ -4,8 +4,8 @@ import ( "context" "strings" - "github.com/aws/aws-sdk-go/aws/session" - "github.com/aws/aws-sdk-go/service/s3" + awsConfig "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/service/s3" "github.com/segmentio/cli" dig "github.com/segmentio/data-digger/pkg/digger" log "github.com/sirupsen/logrus" @@ -38,8 +38,11 @@ func S3Cmd(ctx context.Context) cli.Function { log.Fatalf("Error creating processors: %+v", err) } - sess := session.Must(session.NewSession()) - s3Client := s3.New(sess) + cfg, err := awsConfig.LoadDefaultConfig(ctx) + if err != nil { + log.Fatalf("Unable to load AWS config: %v", err) + } + s3Client := s3.NewFromConfig(cfg) digger := &dig.Digger{ SourceConsumer: &dig.S3Consumer{ diff --git a/go.mod b/go.mod index b526c39..0a80857 100644 --- a/go.mod +++ b/go.mod @@ -3,7 +3,10 @@ module github.com/segmentio/data-digger go 1.24.4 require ( - github.com/aws/aws-sdk-go v1.55.7 + github.com/aws/aws-sdk-go-v2 v1.32.6 + github.com/aws/aws-sdk-go-v2/config v1.28.6 + github.com/aws/aws-sdk-go-v2/credentials v1.17.47 + github.com/aws/aws-sdk-go-v2/service/s3 v1.71.0 github.com/briandowns/spinner v1.23.2 github.com/gogo/protobuf v1.3.2 github.com/gosuri/uilive v0.0.4 @@ -18,9 +21,22 @@ require ( ) require ( + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.7 // indirect + github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.21 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.25 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.25 // indirect + github.com/aws/aws-sdk-go-v2/internal/ini v1.8.1 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.25 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.4.6 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.6 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.6 // indirect + github.com/aws/aws-sdk-go-v2/service/sso v1.24.7 // indirect + github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.6 // indirect + github.com/aws/aws-sdk-go-v2/service/sts v1.33.2 // indirect + github.com/aws/smithy-go v1.22.1 // indirect github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect github.com/fatih/color v1.18.0 // indirect - github.com/jmespath/go-jmespath v0.4.0 // indirect github.com/klauspost/compress v1.18.0 // indirect github.com/mattn/go-colorable v0.1.14 // indirect github.com/mattn/go-isatty v0.0.20 // indirect diff --git a/go.sum b/go.sum index a73763a..413fb65 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,39 @@ -github.com/aws/aws-sdk-go v1.55.7 h1:UJrkFq7es5CShfBwlWAC8DA077vp8PyVbQd3lqLiztE= -github.com/aws/aws-sdk-go v1.55.7/go.mod h1:eRwEWoyTWFMVYVQzKMNHWP5/RV4xIUGMQfXQHfHkpNU= +github.com/aws/aws-sdk-go-v2 v1.32.6 h1:7BokKRgRPuGmKkFMhEg/jSul+tB9VvXhcViILtfG8b4= +github.com/aws/aws-sdk-go-v2 v1.32.6/go.mod h1:P5WJBrYqqbWVaOxgH0X/FYYD47/nooaPOZPlQdmiN2U= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.7 h1:lL7IfaFzngfx0ZwUGOZdsFFnQ5uLvR0hWqqhyE7Q9M8= +github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.7/go.mod h1:QraP0UcVlQJsmHfioCrveWOC1nbiWUl3ej08h4mXWoc= +github.com/aws/aws-sdk-go-v2/config v1.28.6 h1:D89IKtGrs/I3QXOLNTH93NJYtDhm8SYa9Q5CsPShmyo= +github.com/aws/aws-sdk-go-v2/config v1.28.6/go.mod h1:GDzxJ5wyyFSCoLkS+UhGB0dArhb9mI+Co4dHtoTxbko= +github.com/aws/aws-sdk-go-v2/credentials v1.17.47 h1:48bA+3/fCdi2yAwVt+3COvmatZ6jUDNkDTIsqDiMUdw= +github.com/aws/aws-sdk-go-v2/credentials v1.17.47/go.mod h1:+KdckOejLW3Ks3b0E3b5rHsr2f9yuORBum0WPnE5o5w= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.21 h1:AmoU1pziydclFT/xRV+xXE/Vb8fttJCLRPv8oAkprc0= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.21/go.mod h1:AjUdLYe4Tgs6kpH4Bv7uMZo7pottoyHMn4eTcIcneaY= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.25 h1:s/fF4+yDQDoElYhfIVvSNyeCydfbuTKzhxSXDXCPasU= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.25/go.mod h1:IgPfDv5jqFIzQSNbUEMoitNooSMXjRSDkhXv8jiROvU= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.25 h1:ZntTCl5EsYnhN/IygQEUugpdwbhdkom9uHcbCftiGgA= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.25/go.mod h1:DBdPrgeocww+CSl1C8cEV8PN1mHMBhuCDLpXezyvWkE= +github.com/aws/aws-sdk-go-v2/internal/ini v1.8.1 h1:VaRN3TlFdd6KxX1x3ILT5ynH6HvKgqdiXoTxAF4HQcQ= +github.com/aws/aws-sdk-go-v2/internal/ini v1.8.1/go.mod h1:FbtygfRFze9usAadmnGJNc8KsP346kEe+y2/oyhGAGc= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.25 h1:r67ps7oHCYnflpgDy2LZU0MAQtQbYIOqNNnqGO6xQkE= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.25/go.mod h1:GrGY+Q4fIokYLtjCVB/aFfCVL6hhGUFl8inD18fDalE= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.1 h1:iXtILhvDxB6kPvEXgsDhGaZCSC6LQET5ZHSdJozeI0Y= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.12.1/go.mod h1:9nu0fVANtYiAePIBh2/pFUSwtJ402hLnp854CNoDOeE= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.4.6 h1:HCpPsWqmYQieU7SS6E9HXfdAMSud0pteVXieJmcpIRI= +github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.4.6/go.mod h1:ngUiVRCco++u+soRRVBIvBZxSMMvOVMXA4PJ36JLfSw= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.6 h1:50+XsN70RS7dwJ2CkVNXzj7U2L1HKP8nqTd3XWEXBN4= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.12.6/go.mod h1:WqgLmwY7so32kG01zD8CPTJWVWM+TzJoOVHwTg4aPug= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.6 h1:BbGDtTi0T1DYlmjBiCr/le3wzhA37O8QTC5/Ab8+EXk= +github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.18.6/go.mod h1:hLMJt7Q8ePgViKupeymbqI0la+t9/iYFBjxQCFwuAwI= +github.com/aws/aws-sdk-go-v2/service/s3 v1.71.0 h1:nyuzXooUNJexRT0Oy0UQY6AhOzxPxhtt4DcBIHyCnmw= +github.com/aws/aws-sdk-go-v2/service/s3 v1.71.0/go.mod h1:sT/iQz8JK3u/5gZkT+Hmr7GzVZehUMkRZpOaAwYXeGY= +github.com/aws/aws-sdk-go-v2/service/sso v1.24.7 h1:rLnYAfXQ3YAccocshIH5mzNNwZBkBo+bP6EhIxak6Hw= +github.com/aws/aws-sdk-go-v2/service/sso v1.24.7/go.mod h1:ZHtuQJ6t9A/+YDuxOLnbryAmITtr8UysSny3qcyvJTc= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.6 h1:JnhTZR3PiYDNKlXy50/pNeix9aGMo6lLpXwJ1mw8MD4= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.28.6/go.mod h1:URronUEGfXZN1VpdktPSD1EkAL9mfrV+2F4sjH38qOY= +github.com/aws/aws-sdk-go-v2/service/sts v1.33.2 h1:s4074ZO1Hk8qv65GqNXqDjmkf4HSQqJukaLuuW0TpDA= +github.com/aws/aws-sdk-go-v2/service/sts v1.33.2/go.mod h1:mVggCnIWoM09jP71Wh+ea7+5gAp53q+49wDFs1SW5z8= +github.com/aws/smithy-go v1.22.1 h1:/HPHZQ0g7f4eUeK6HKglFz8uwVfZKgoI25rb/J+dnro= +github.com/aws/smithy-go v1.22.1/go.mod h1:irrKGvNn1InZwb2d7fkIRNucdfwR8R+Ts3wxYa/cJHg= github.com/briandowns/spinner v1.23.2 h1:Zc6ecUnI+YzLmJniCfDNaMbW0Wid1d5+qcTq4L2FW8w= github.com/briandowns/spinner v1.23.2/go.mod h1:LaZeM4wm2Ywy6vO571mvhQNRcWfRUnXOs0RcKV0wYKM= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -16,10 +50,6 @@ github.com/gosuri/uilive v0.0.4 h1:hUEBpQDj8D8jXgtCdBu7sWsy5sbW/5GhuO8KBwJ2jyY= github.com/gosuri/uilive v0.0.4/go.mod h1:V/epo5LjjlDE5RJUcqx8dbw+zc93y5Ya3yg8tfZ74VI= github.com/hpcloud/tail v1.0.0 h1:nfCOvKYfkgYP8hkirhJocXT2+zOD8yUNjXaWfTlyFKI= github.com/hpcloud/tail v1.0.0/go.mod h1:ab1qPbhIpdTxEkNHXyeSf5vhxWSCs/tWer42PpOxQnU= -github.com/jmespath/go-jmespath v0.4.0 h1:BEgLn5cpjn8UN1mAw4NjwDrS35OdebyEtFe+9YPoQUg= -github.com/jmespath/go-jmespath v0.4.0/go.mod h1:T8mJZnbsbmF+m6zOOFylbeCJqk5+pHWvzYPziyZiYoo= -github.com/jmespath/go-jmespath/internal/testify v1.5.1 h1:shLQSRRSCCPj3f2gpwzGwWFoC7ycTf1rcQZHOlsJ6N8= -github.com/jmespath/go-jmespath/internal/testify v1.5.1/go.mod h1:L3OGu8Wl2/fWfCI6z80xFu9LTZmf1ZRjMHUOPmWr69U= github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= github.com/klauspost/compress v1.15.9/go.mod h1:PhcZ0MbTNciWF3rruxRgKxI5NkcHHrHUDtV4Yw2GlzU= @@ -161,7 +191,6 @@ gopkg.in/fsnotify.v1 v1.4.7/go.mod h1:Tz8NjZHkW78fSQdbUxIjBTcgA1z1m8ZHf0WmKUhAMy gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7 h1:uRGJdciOHaEIrze2W8Q3AKkepLTh2hOroT7a+7czfdQ= gopkg.in/tomb.v1 v1.0.0-20141024135613-dd632973f1e7/go.mod h1:dt/ZhP58zS4L8KSrWDmTeBkI65Dw0HsyUHuEVlX15mw= gopkg.in/yaml.v2 v2.2.1/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v2 v2.2.8/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= gopkg.in/yaml.v2 v2.3.0 h1:clyUAQHOM3G0M3f5vQj7LuJrETvjVot3Z5el9nffUtU= gopkg.in/yaml.v2 v2.3.0/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/pkg/digger/s3.go b/pkg/digger/s3.go index f44248c..21d1b72 100644 --- a/pkg/digger/s3.go +++ b/pkg/digger/s3.go @@ -7,8 +7,9 @@ import ( "fmt" "strings" - "github.com/aws/aws-sdk-go/aws" - "github.com/aws/aws-sdk-go/service/s3" + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/aws/aws-sdk-go-v2/service/s3/types" "github.com/segmentio/kafka-go" log "github.com/sirupsen/logrus" ) @@ -16,7 +17,7 @@ import ( // S3Consumer is a Consumer implementation that reads from one or more prefixes in an S3 // bucket. type S3Consumer struct { - S3Client *s3.S3 + S3Client *s3.Client Bucket string Prefixes []string NumWorkers int @@ -25,7 +26,7 @@ type S3Consumer struct { var _ Consumer = (*S3Consumer)(nil) type s3ObjTask struct { - objInfo *s3.Object + objInfo types.Object index int } @@ -64,34 +65,31 @@ func (s *S3Consumer) processPrefixes( prefixesRead := 0 for _, prefix := range s.Prefixes { - err := s.S3Client.ListObjectsPagesWithContext( - ctx, - &s3.ListObjectsInput{ - Bucket: aws.String(s.Bucket), - Prefix: aws.String(prefix), - }, - func(output *s3.ListObjectsOutput, hasMore bool) bool { - for _, objInfo := range output.Contents { - subTask := s3ObjTask{ - objInfo: objInfo, - index: keysRead, - } - select { - case objectChan <- subTask: - case <-ctx.Done(): - return false - } - keysRead++ - } + paginator := s3.NewListObjectsV2Paginator(s.S3Client, &s3.ListObjectsV2Input{ + Bucket: aws.String(s.Bucket), + Prefix: aws.String(prefix), + }) - return true - }, - ) + for paginator.HasMorePages() { + output, err := paginator.NextPage(ctx) + if err != nil { + return err + } - prefixesRead++ - if err != nil { - return err + for _, objInfo := range output.Contents { + subTask := s3ObjTask{ + objInfo: objInfo, + index: keysRead, + } + select { + case objectChan <- subTask: + case <-ctx.Done(): + return ctx.Err() + } + keysRead++ + } } + prefixesRead++ } return nil @@ -113,7 +111,7 @@ func (s *S3Consumer) runSubTasks( if err != nil { return fmt.Errorf( "Error processing key %s: %+v", - aws.StringValue(subTask.objInfo.Key), + aws.ToString(subTask.objInfo.Key), err, ) } @@ -126,18 +124,18 @@ func (s *S3Consumer) runSubTasks( func (s *S3Consumer) processKey( ctx context.Context, messageChan chan message, - objInfo *s3.Object, + objInfo types.Object, index int, ) error { - log.Debugf("Processing key %s", aws.StringValue(objInfo.Key)) + log.Debugf("Processing key %s", aws.ToString(objInfo.Key)) var contentEncoding *string - if strings.HasSuffix(aws.StringValue(objInfo.Key), ".gz") { + if strings.HasSuffix(aws.ToString(objInfo.Key), ".gz") { // Assume gzip encoding (which might not actually be set in the object in S3) contentEncoding = aws.String("gzip") } - obj, err := s.S3Client.GetObjectWithContext( + obj, err := s.S3Client.GetObject( ctx, &s3.GetObjectInput{ Bucket: aws.String(s.Bucket), @@ -175,8 +173,8 @@ func (s *S3Consumer) processKey( messageChan <- message{ msg: kafka.Message{ Partition: index, - Time: aws.TimeValue(objInfo.LastModified), - Key: []byte(aws.StringValue(objInfo.Key)), + Time: aws.ToTime(objInfo.LastModified), + Key: []byte(aws.ToString(objInfo.Key)), Offset: offset, Value: copiedContents, }, diff --git a/pkg/digger/s3_test.go b/pkg/digger/s3_test.go index 2d732a6..fd909db 100644 --- a/pkg/digger/s3_test.go +++ b/pkg/digger/s3_test.go @@ -8,17 +8,16 @@ import ( "testing" "time" - "github.com/aws/aws-sdk-go/aws" - "github.com/aws/aws-sdk-go/aws/credentials" - "github.com/aws/aws-sdk-go/aws/session" - "github.com/aws/aws-sdk-go/service/s3" + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/credentials" + "github.com/aws/aws-sdk-go-v2/service/s3" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) func TestS3Consumer(t *testing.T) { ctx := context.Background() - sess := session.Must(session.NewSession()) var s3Endpoint string @@ -29,20 +28,28 @@ func TestS3Consumer(t *testing.T) { s3Endpoint = "http://localhost:4572" } - s3Client := s3.New( - sess, - &aws.Config{ + cfg, err := config.LoadDefaultConfig(ctx, + config.WithCredentialsProvider( // These need to be set, but they can be anything since localstack // doesn't do any checking - Credentials: credentials.NewStaticCredentials("test", "test", "test"), - - Endpoint: aws.String(s3Endpoint), - Region: aws.String("us-west-2"), - DisableSSL: aws.Bool(true), - S3ForcePathStyle: aws.Bool(true), - }, + credentials.NewStaticCredentialsProvider("test", "test", "test"), + ), + config.WithRegion("us-west-2"), + config.WithEndpointResolverWithOptions( + aws.EndpointResolverWithOptionsFunc(func(service, region string, options ...interface{}) (aws.Endpoint, error) { + return aws.Endpoint{ + URL: s3Endpoint, + HostnameImmutable: true, + SigningRegion: "us-west-2", + }, nil + }), + ), ) + s3Client := s3.NewFromConfig(cfg, func(o *s3.Options) { + o.UsePathStyle = true + }) + testBucket := createBucket(ctx, t, s3Client) time.Sleep(100 * time.Millisecond) @@ -57,7 +64,7 @@ func TestS3Consumer(t *testing.T) { Prefixes: []string{"test-prefix1", "test-prefix2"}, NumWorkers: 1, } - err := consumer.Run(ctx, messageChan) + err = consumer.Run(ctx, messageChan) require.NoError(t, err) require.Equal(t, 5, len(messageChan)) @@ -74,10 +81,10 @@ func TestS3Consumer(t *testing.T) { assert.Equal(t, []byte("value2"), message2.msg.Value) } -func createBucket(ctx context.Context, t *testing.T, s3Client *s3.S3) string { +func createBucket(ctx context.Context, t *testing.T, s3Client *s3.Client) string { bucketName := fmt.Sprintf("test-bucket-%d", time.Now().UnixNano()) - _, err := s3Client.CreateBucketWithContext( + _, err := s3Client.CreateBucket( ctx, &s3.CreateBucketInput{ Bucket: aws.String(bucketName), @@ -90,19 +97,19 @@ func createBucket(ctx context.Context, t *testing.T, s3Client *s3.S3) string { func writeKey( ctx context.Context, t *testing.T, - s3Client *s3.S3, + s3Client *s3.Client, bucket string, key string, value string, ) { body := bytes.NewBufferString(value) - _, err := s3Client.PutObjectWithContext( + _, err := s3Client.PutObject( ctx, &s3.PutObjectInput{ Bucket: aws.String(bucket), Key: aws.String(key), - Body: aws.ReadSeekCloser(body), + Body: body, }, ) require.NoError(t, err)