Преглед на файлове

feat(aws): add Athena query result reuse configuration (#3953)

Signed-off-by: Sridhar Vemula <thewarrior.316@gmail.com>
Sridhar Vemula преди 1 месец
родител
ревизия
e7e93745df
променени са 3 файла, в които са добавени 81 реда и са изтрити 8 реда
  1. 25 8
      pkg/cloud/aws/athenaconfiguration.go
  2. 47 0
      pkg/cloud/aws/athenaconfiguration_test.go
  3. 9 0
      pkg/cloud/aws/athenaquerier.go

+ 25 - 8
pkg/cloud/aws/athenaconfiguration.go

@@ -18,6 +18,10 @@ type AthenaConfiguration struct {
 	Workgroup  string     `json:"workgroup"`
 	Account    string     `json:"account"`
 	Authorizer Authorizer `json:"authorizer"`
+	// ResultReuseMaxAgeMinutes enables Athena query result reuse for this many minutes (0 = disabled).
+	// When set, identical queries within the TTL window return cached results instead of re-scanning,
+	// reducing cost for multi-cluster deployments where many clusters run the same CUR query.
+	ResultReuseMaxAgeMinutes int32 `json:"resultReuseMaxAgeMinutes,omitempty"`
 }
 
 func (ac *AthenaConfiguration) Validate() error {
@@ -103,19 +107,24 @@ func (ac *AthenaConfiguration) Equals(config cloud.Config) bool {
 		return false
 	}
 
+	if ac.ResultReuseMaxAgeMinutes != thatConfig.ResultReuseMaxAgeMinutes {
+		return false
+	}
+
 	return true
 }
 
 func (ac *AthenaConfiguration) Sanitize() cloud.Config {
 	return &AthenaConfiguration{
-		Bucket:     ac.Bucket,
-		Region:     ac.Region,
-		Database:   ac.Database,
-		Catalog:    ac.Catalog,
-		Table:      ac.Table,
-		Workgroup:  ac.Workgroup,
-		Account:    ac.Account,
-		Authorizer: ac.Authorizer.Sanitize().(Authorizer),
+		Bucket:                   ac.Bucket,
+		Region:                   ac.Region,
+		Database:                 ac.Database,
+		Catalog:                  ac.Catalog,
+		Table:                    ac.Table,
+		Workgroup:                ac.Workgroup,
+		Account:                  ac.Account,
+		Authorizer:               ac.Authorizer.Sanitize().(Authorizer),
+		ResultReuseMaxAgeMinutes: ac.ResultReuseMaxAgeMinutes,
 	}
 }
 
@@ -190,6 +199,14 @@ func (ac *AthenaConfiguration) UnmarshalJSON(b []byte) error {
 	}
 	ac.Authorizer = authorizer
 
+	if _, ok := fmap["resultReuseMaxAgeMinutes"]; ok {
+		resultReuseMaxAgeMinutes, err := cloud.GetInterfaceValue[float64](fmap, "resultReuseMaxAgeMinutes")
+		if err != nil {
+			return fmt.Errorf("AthenaConfiguration: UnmarshalJSON: %w", err)
+		}
+		ac.ResultReuseMaxAgeMinutes = int32(resultReuseMaxAgeMinutes)
+	}
+
 	return nil
 }
 

+ 47 - 0
pkg/cloud/aws/athenaconfiguration_test.go

@@ -515,6 +515,37 @@ func TestAthenaConfiguration_Equals(t *testing.T) {
 			},
 			expected: false,
 		},
+		"different resultReuseMaxAgeMinutes": {
+			left: AthenaConfiguration{
+				Bucket:                   "bucket",
+				Region:                   "region",
+				Database:                 "database",
+				Catalog:                  "catalog",
+				Table:                    "table",
+				Workgroup:                "workgroup",
+				Account:                  "account",
+				ResultReuseMaxAgeMinutes: 60,
+				Authorizer: &AccessKey{
+					ID:     "id",
+					Secret: "secret",
+				},
+			},
+			right: &AthenaConfiguration{
+				Bucket:                   "bucket",
+				Region:                   "region",
+				Database:                 "database",
+				Catalog:                  "catalog",
+				Table:                    "table",
+				Workgroup:                "workgroup",
+				Account:                  "account",
+				ResultReuseMaxAgeMinutes: 120,
+				Authorizer: &AccessKey{
+					ID:     "id",
+					Secret: "secret",
+				},
+			},
+			expected: false,
+		},
 		"different config": {
 			left: AthenaConfiguration{
 				Bucket:    "bucket",
@@ -581,6 +612,22 @@ func TestAthenaConfiguration_JSON(t *testing.T) {
 				Authorizer: &ServiceAccount{},
 			},
 		},
+		"ResultReuseMaxAgeMinutes": {
+			config: AthenaConfiguration{
+				Bucket:                   "bucket",
+				Region:                   "region",
+				Database:                 "database",
+				Catalog:                  "catalog",
+				Table:                    "table",
+				Workgroup:                "workgroup",
+				Account:                  "account",
+				ResultReuseMaxAgeMinutes: 60,
+				Authorizer: &AccessKey{
+					ID:     "id",
+					Secret: "secret",
+				},
+			},
+		},
 		"AssumeRole with AccessKey": {
 			config: AthenaConfiguration{
 				Bucket:    "bucket",

+ 9 - 0
pkg/cloud/aws/athenaquerier.go

@@ -113,6 +113,15 @@ func (aq *AthenaQuerier) queryAthenaPaginated(ctx context.Context, query string,
 		startQueryExecutionInput.WorkGroup = aws.String(aq.Workgroup)
 	}
 
+	if aq.ResultReuseMaxAgeMinutes > 0 {
+		startQueryExecutionInput.ResultReuseConfiguration = &types.ResultReuseConfiguration{
+			ResultReuseByAgeConfiguration: &types.ResultReuseByAgeConfiguration{
+				Enabled:         true,
+				MaxAgeInMinutes: aws.Int32(aq.ResultReuseMaxAgeMinutes),
+			},
+		}
+	}
+
 	// Create Athena Client
 	cli, err := aq.GetAthenaClient()
 	if err != nil {