Skip to content

Commit 796ffad

Browse files
author
Mike Paßberg
authored
Add support for AWS SDK v2 and DynamoDB to CSV export (#8)
1 parent e23bdf3 commit 796ffad

4 files changed

Lines changed: 412 additions & 0 deletions

File tree

pkg/csvexport/csv_v2.go

Lines changed: 191 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,191 @@
1+
package csvexport
2+
3+
import (
4+
"bytes"
5+
"context"
6+
"encoding/csv"
7+
"encoding/json"
8+
"fmt"
9+
"sort"
10+
"strconv"
11+
12+
"github.com/aws/aws-sdk-go-v2/aws"
13+
"github.com/aws/aws-sdk-go-v2/feature/dynamodb/attributevalue"
14+
"github.com/aws/aws-sdk-go-v2/service/dynamodb"
15+
"github.com/aws/aws-sdk-go-v2/service/dynamodb/types"
16+
"github.com/aws/aws-sdk-go-v2/service/s3"
17+
)
18+
19+
type StorageV2 interface {
20+
Scan(ctx context.Context, opt ScanOption, startKey map[string]types.AttributeValue) ([]map[string]interface{}, map[string]types.AttributeValue, error)
21+
}
22+
23+
type DynoStorageV2 struct {
24+
DDB *dynamodb.Client
25+
}
26+
27+
func DynamoToCSVV2(db StorageV2, ctx context.Context, scanOpt ScanOption, opts ...Option) ([]byte, error) {
28+
var b bytes.Buffer
29+
30+
w := csv.NewWriter(&b)
31+
32+
defer w.Flush()
33+
34+
var startKey map[string]types.AttributeValue
35+
36+
var csvExp CSVExporter
37+
38+
var keyOrder []string
39+
var header []string
40+
41+
for _, opt := range opts {
42+
opt(&csvExp)
43+
}
44+
45+
count := 0
46+
47+
for {
48+
resp, sk, err := db.Scan(ctx, scanOpt, startKey)
49+
if err != nil {
50+
return nil, fmt.Errorf("scan failed: %w", err)
51+
}
52+
53+
for _, attr := range resp {
54+
if count == 0 {
55+
for _, v := range csvExp.cols {
56+
keyOrder = append(keyOrder, v.Name)
57+
58+
headerName := v.Name
59+
60+
if v.TargetName != "" {
61+
headerName = v.TargetName
62+
}
63+
64+
header = append(header, headerName)
65+
}
66+
67+
if csvExp.cols == nil {
68+
for k := range attr {
69+
keyOrder = append(keyOrder, k)
70+
}
71+
72+
sort.Strings(keyOrder)
73+
74+
header = keyOrder
75+
}
76+
77+
_ = w.Write(header)
78+
}
79+
80+
record := make([]string, 0, len(keyOrder))
81+
82+
for i, k := range keyOrder {
83+
value := attr[k]
84+
85+
// Empty Value of column?
86+
if len(csvExp.cols) > 0 && csvExp.cols[i].OverwriteValue {
87+
value = csvExp.cols[i].OverwriteWithValue
88+
}
89+
90+
// Custom function?
91+
valueFn, valueFnCol, ok := csvExp.cols.ValueFunc(k)
92+
93+
if ok {
94+
if valueFnCol != "" {
95+
value = attr[valueFnCol]
96+
}
97+
newVal, err := valueFn(value)
98+
if err != nil {
99+
return nil, fmt.Errorf("failed to process custom valueFunction on column %s: %w", valueFnCol, err)
100+
}
101+
record = append(record, newVal)
102+
continue
103+
}
104+
105+
switch val := value.(type) {
106+
case float64:
107+
// protect exponential notation layout
108+
record = append(record, strconv.FormatFloat(val, 'f', -1, 64))
109+
case string:
110+
record = append(record, removeNewLines(val))
111+
default:
112+
js, err := json.Marshal(value)
113+
if err != nil {
114+
return nil, fmt.Errorf("failed to marshal value: %w", err)
115+
}
116+
117+
record = append(record, string(js))
118+
}
119+
}
120+
121+
_ = w.Write(record)
122+
count++
123+
}
124+
125+
startKey = sk
126+
if len(startKey) == 0 {
127+
break
128+
}
129+
}
130+
131+
w.Flush()
132+
133+
return b.Bytes(), nil
134+
}
135+
136+
func (db DynoStorageV2) Scan(ctx context.Context, opt ScanOption, startKey map[string]types.AttributeValue) ([]map[string]interface{}, map[string]types.AttributeValue, error) {
137+
var expressionAttributeValues map[string]types.AttributeValue
138+
if opt.ExpressionAttrValues != "" {
139+
if err := json.Unmarshal([]byte(opt.ExpressionAttrValues), &expressionAttributeValues); err != nil {
140+
return nil, nil, fmt.Errorf("expression attribute values invalid: %w", err)
141+
}
142+
}
143+
144+
var expressionAttributeNames map[string]string
145+
if opt.ExpressionAttrNames != "" {
146+
expressionAttributeNames = make(map[string]string)
147+
if err := json.Unmarshal([]byte(opt.ExpressionAttrNames), &expressionAttributeNames); err != nil {
148+
return nil, nil, fmt.Errorf("expression attribute names invalid: %w", err)
149+
}
150+
}
151+
152+
var filterExpression *string
153+
if opt.FilterExpression != "" {
154+
filterExpression = aws.String(opt.FilterExpression)
155+
}
156+
157+
out, err := db.DDB.Scan(ctx, &dynamodb.ScanInput{
158+
TableName: aws.String(opt.TableName),
159+
ExclusiveStartKey: startKey,
160+
ExpressionAttributeNames: expressionAttributeNames,
161+
FilterExpression: filterExpression,
162+
ExpressionAttributeValues: expressionAttributeValues,
163+
})
164+
if err != nil {
165+
return nil, nil, fmt.Errorf("db.Scan: %w", err)
166+
}
167+
168+
var resp []map[string]interface{}
169+
if err := attributevalue.UnmarshalListOfMaps(out.Items, &resp); err != nil {
170+
return nil, nil, fmt.Errorf("dynamodb unmarshal list of maps: %w", err)
171+
}
172+
173+
return resp, out.LastEvaluatedKey, nil
174+
}
175+
176+
type S3PutClient interface {
177+
PutObject(ctx context.Context, params *s3.PutObjectInput, optFns ...func(*s3.Options)) (*s3.PutObjectOutput, error)
178+
}
179+
180+
func UploadToS3V2(ctx context.Context, b []byte, client S3PutClient, bucket, path, fname string) error {
181+
_, err := client.PutObject(ctx, &s3.PutObjectInput{
182+
Body: bytes.NewReader(b),
183+
Bucket: aws.String(bucket),
184+
Key: aws.String(path + fname),
185+
})
186+
if err != nil {
187+
return fmt.Errorf("failed to PutObject %s to S3: %w", fname, err)
188+
}
189+
190+
return nil
191+
}

pkg/csvexport/csv_v2_test.go

Lines changed: 176 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,176 @@
1+
package csvexport_test
2+
3+
import (
4+
"context"
5+
"fmt"
6+
"github.com/aws/aws-sdk-go-v2/service/dynamodb/types"
7+
"testing"
8+
9+
"github.com/spring-media/curation-pkgs-public/pkg/csvexport"
10+
"github.com/stretchr/testify/assert"
11+
)
12+
13+
var dynamoMockRespV2 = []map[string]interface{}{
14+
{
15+
"ArticleLastUpdated": "2022-04-27T08:36:48.386Z",
16+
"Block": "Meldungen1",
17+
"Categories": []interface{}{
18+
map[string]interface{}{
19+
"name": "money_business1",
20+
"score": 0.958231,
21+
},
22+
},
23+
"EntriesAveragedFromHome": map[string]interface{}{
24+
"15": 1.333333,
25+
"180": 2.861111,
26+
"30": 3.166667,
27+
"45": 4.888889,
28+
"5": 5.000000,
29+
},
30+
"IsPremium": false,
31+
"IsSponsored": false,
32+
"Meta": map[string]interface{}{
33+
"Department": "wirtschaft",
34+
"FromInvestigativ": false,
35+
"FromNewsteam": true,
36+
"HasVideo": false,
37+
},
38+
"PerformanceLastUpdated": "2022-04-27T13:10:30Z",
39+
},
40+
{
41+
"ArticleLastUpdated": "2022-04-28T08:36:48.386Z",
42+
"Block": "Meldungen2",
43+
"Categories": []interface{}{
44+
map[string]interface{}{
45+
"name": "money_business2",
46+
"score": 0.99,
47+
},
48+
},
49+
"EntriesAveragedFromHome": map[string]interface{}{
50+
"15": 6.333333,
51+
"180": 7.861111,
52+
"30": 8.166667,
53+
"45": 9.888889,
54+
"5": 10.000000,
55+
},
56+
"IsPremium": true,
57+
"IsSponsored": true,
58+
"Meta": map[string]interface{}{
59+
"Department": "wirtschaft2",
60+
"FromInvestigativ": true,
61+
"FromNewsteam": false,
62+
"HasVideo": true,
63+
},
64+
"PerformanceLastUpdated": "2022-04-28T13:10:30Z",
65+
},
66+
}
67+
68+
type mockScanV2 struct {
69+
resp []map[string]interface{}
70+
}
71+
72+
func (d mockScanV2) Scan(ctx context.Context, opt csvexport.ScanOption, startKey map[string]types.AttributeValue) ([]map[string]interface{}, map[string]types.AttributeValue, error) {
73+
return d.resp, map[string]types.AttributeValue{}, nil
74+
}
75+
76+
func TestDynamoToCSVV2(t *testing.T) {
77+
t.Parallel()
78+
option := csvexport.ScanOption{}
79+
80+
db := mockScanV2{resp: dynamoMockRespV2}
81+
b, err := csvexport.DynamoToCSVV2(db, context.Background(), option)
82+
83+
assert.Nil(t, err)
84+
85+
expectedCSV := `ArticleLastUpdated,Block,Categories,EntriesAveragedFromHome,IsPremium,IsSponsored,Meta,PerformanceLastUpdated
86+
2022-04-27T08:36:48.386Z,Meldungen1,"[{""name"":""money_business1"",""score"":0.958231}]","{""15"":1.333333,""180"":2.861111,""30"":3.166667,""45"":4.888889,""5"":5}",false,false,"{""Department"":""wirtschaft"",""FromInvestigativ"":false,""FromNewsteam"":true,""HasVideo"":false}",2022-04-27T13:10:30Z
87+
2022-04-28T08:36:48.386Z,Meldungen2,"[{""name"":""money_business2"",""score"":0.99}]","{""15"":6.333333,""180"":7.861111,""30"":8.166667,""45"":9.888889,""5"":10}",true,true,"{""Department"":""wirtschaft2"",""FromInvestigativ"":true,""FromNewsteam"":false,""HasVideo"":true}",2022-04-28T13:10:30Z
88+
`
89+
assert.Equal(t, expectedCSV, string(b))
90+
}
91+
92+
func TestDynamoToCSVWithColsV2(t *testing.T) {
93+
t.Parallel()
94+
db := mockScanV2{resp: dynamoMockRespV2}
95+
cols := csvexport.Columns{
96+
csvexport.Column{Name: "ArticleLastUpdated"},
97+
csvexport.Column{Name: "Block", OverwriteValue: true, OverwriteWithValue: ""},
98+
csvexport.Column{Name: "Categories", OverwriteValue: true, OverwriteWithValue: nil},
99+
csvexport.Column{Name: "EntriesAveragedFromHome", OverwriteValue: true, OverwriteWithValue: 3},
100+
csvexport.Column{Name: "IsPremium", OverwriteValue: true, OverwriteWithValue: 1.33},
101+
csvexport.Column{Name: "IsSponsored", OverwriteValue: true, OverwriteWithValue: false},
102+
csvexport.Column{Name: "PerformanceLastUpdated"},
103+
}
104+
b, err := csvexport.DynamoToCSVV2(db, context.Background(), csvexport.ScanOption{}, csvexport.WithColumns(cols))
105+
106+
assert.Nil(t, err)
107+
108+
expectedCSV := `ArticleLastUpdated,Block,Categories,EntriesAveragedFromHome,IsPremium,IsSponsored,PerformanceLastUpdated
109+
2022-04-27T08:36:48.386Z,,null,3,1.33,false,2022-04-27T13:10:30Z
110+
2022-04-28T08:36:48.386Z,,null,3,1.33,false,2022-04-28T13:10:30Z
111+
`
112+
assert.Equal(t, expectedCSV, string(b))
113+
}
114+
115+
func TestDynamoToCSVWithColsTargetNameV2(t *testing.T) {
116+
t.Parallel()
117+
db := mockScanV2{resp: dynamoMockRespV2}
118+
cols := csvexport.Columns{
119+
csvexport.Column{Name: "ArticleLastUpdated"},
120+
csvexport.Column{Name: "PerformanceLastUpdated", TargetName: "PerformanceUpdatedLast"},
121+
}
122+
b, err := csvexport.DynamoToCSVV2(db, context.Background(), csvexport.ScanOption{}, csvexport.WithColumns(cols))
123+
124+
assert.Nil(t, err)
125+
126+
expectedCSV := `ArticleLastUpdated,PerformanceUpdatedLast
127+
2022-04-27T08:36:48.386Z,2022-04-27T13:10:30Z
128+
2022-04-28T08:36:48.386Z,2022-04-28T13:10:30Z
129+
`
130+
assert.Equal(t, expectedCSV, string(b))
131+
}
132+
133+
func TestDynamoToCSVWithColsValueFuncV2(t *testing.T) {
134+
t.Parallel()
135+
db := mockScanV2{resp: dynamoMockRespV2}
136+
137+
valueFn := func(val interface{}) (string, error) {
138+
v, ok := val.(string)
139+
if !ok {
140+
return "", fmt.Errorf("failed to cast to string")
141+
}
142+
143+
if v == "Meldungen1" {
144+
return "Meldungen123", nil
145+
}
146+
147+
return v, nil
148+
}
149+
150+
valueFnNew := func(val interface{}) (string, error) {
151+
v, ok := val.(string)
152+
if !ok {
153+
return "", fmt.Errorf("failed to cast to string")
154+
}
155+
156+
if v == "Meldungen2" {
157+
return "Meldungen999", nil
158+
}
159+
160+
return v, nil
161+
}
162+
cols := csvexport.Columns{
163+
csvexport.Column{Name: "ArticleLastUpdated"},
164+
csvexport.Column{Name: "Block", ValueFunc: valueFn},
165+
csvexport.Column{Name: "NewBlock", ValueFunc: valueFnNew, ValueFuncCol: "Block"},
166+
}
167+
b, err := csvexport.DynamoToCSVV2(db, context.Background(), csvexport.ScanOption{}, csvexport.WithColumns(cols))
168+
169+
assert.Nil(t, err)
170+
171+
expectedCSV := `ArticleLastUpdated,Block,NewBlock
172+
2022-04-27T08:36:48.386Z,Meldungen123,Meldungen1
173+
2022-04-28T08:36:48.386Z,Meldungen2,Meldungen999
174+
`
175+
assert.Equal(t, expectedCSV, string(b))
176+
}

pkg/csvexport/go.mod

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,10 +4,25 @@ go 1.22
44

55
require (
66
github.com/aws/aws-sdk-go v1.50.32
7+
github.com/aws/aws-sdk-go-v2 v1.37.0
8+
github.com/aws/aws-sdk-go-v2/feature/dynamodb/attributevalue v1.20.0
9+
github.com/aws/aws-sdk-go-v2/service/dynamodb v1.45.0
10+
github.com/aws/aws-sdk-go-v2/service/s3 v1.85.0
711
github.com/stretchr/testify v1.7.1
812
)
913

1014
require (
15+
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.0 // indirect
16+
github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.0 // indirect
17+
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.0 // indirect
18+
github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.0 // indirect
19+
github.com/aws/aws-sdk-go-v2/service/dynamodbstreams v1.27.0 // indirect
20+
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.0 // indirect
21+
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.8.0 // indirect
22+
github.com/aws/aws-sdk-go-v2/service/internal/endpoint-discovery v1.11.0 // indirect
23+
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.0 // indirect
24+
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.0 // indirect
25+
github.com/aws/smithy-go v1.22.5 // indirect
1126
github.com/davecgh/go-spew v1.1.0 // indirect
1227
github.com/jmespath/go-jmespath v0.4.0 // indirect
1328
github.com/pmezard/go-difflib v1.0.0 // indirect

0 commit comments

Comments
 (0)