Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 7 additions & 4 deletions cmd/digger/subcmd/s3.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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{
Expand Down
20 changes: 18 additions & 2 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down
43 changes: 36 additions & 7 deletions go.sum
Original file line number Diff line number Diff line change
@@ -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=
Expand All @@ -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=
Expand Down Expand Up @@ -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=
Expand Down
70 changes: 34 additions & 36 deletions pkg/digger/s3.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,16 +7,17 @@ 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"
)

// 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
Expand All @@ -25,7 +26,7 @@ type S3Consumer struct {
var _ Consumer = (*S3Consumer)(nil)

type s3ObjTask struct {
objInfo *s3.Object
objInfo types.Object
index int
}

Expand Down Expand Up @@ -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
Expand All @@ -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,
)
}
Expand All @@ -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),
Expand Down Expand Up @@ -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,
},
Expand Down
49 changes: 28 additions & 21 deletions pkg/digger/s3_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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)
Expand All @@ -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))
Expand All @@ -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),
Expand All @@ -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)
Expand Down