Skip to content

Commit e902f7f

Browse files
Feat/scala steward upgrades 5 (#72)
* Update xz to 1.10 * Update s3, sts to 2.26.29 * Use RetryStrategy instead of RetryPolicy --------- Co-authored-by: Scala Steward <scala_steward@virtuslab.com>
1 parent 8513235 commit e902f7f

2 files changed

Lines changed: 56 additions & 34 deletions

File tree

kafka-connect-aws-s3/src/main/scala/io/lenses/streamreactor/connect/aws/s3/auth/AwsS3ClientCreator.scala

Lines changed: 54 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -25,10 +25,10 @@ import software.amazon.awssdk.auth.credentials.AwsCredentials
2525
import software.amazon.awssdk.auth.credentials.AwsCredentialsProvider
2626
import software.amazon.awssdk.auth.credentials.DefaultCredentialsProvider
2727
import software.amazon.awssdk.core.client.config.ClientOverrideConfiguration
28-
import software.amazon.awssdk.core.retry.RetryPolicy
29-
import software.amazon.awssdk.core.retry.backoff.FixedDelayBackoffStrategy
3028
import software.amazon.awssdk.http.apache.ApacheHttpClient
3129
import software.amazon.awssdk.regions.Region
30+
import software.amazon.awssdk.retries.api.RetryStrategy
31+
import software.amazon.awssdk.retries.api.internal.backoff.FixedDelayWithoutJitter
3232
import software.amazon.awssdk.services.s3.S3Client
3333
import software.amazon.awssdk.services.s3.S3Configuration
3434

@@ -45,18 +45,6 @@ object AwsS3ClientCreator extends ClientCreator[S3ConnectionConfig, S3Client] {
4545

4646
def make(config: S3ConnectionConfig): Either[Throwable, S3Client] =
4747
for {
48-
retryPolicy <- Try {
49-
RetryPolicy
50-
.builder()
51-
.numRetries(config.httpRetryConfig.getRetryLimit)
52-
.backoffStrategy(
53-
FixedDelayBackoffStrategy.create(Duration.ofMillis(config.httpRetryConfig.getRetryIntervalMillis)),
54-
)
55-
.build()
56-
}.toEither
57-
58-
overrideConfig <- Try(ClientOverrideConfiguration.builder().retryPolicy(retryPolicy).build()).toEither
59-
6048
s3Config <- Try {
6149
S3Configuration
6250
.builder
@@ -71,27 +59,61 @@ object AwsS3ClientCreator extends ClientCreator[S3ConnectionConfig, S3Client] {
7159
config.connectionPoolConfig.foreach(t => apacheHttpClientBuilder.maxConnections(t.maxConnections))
7260
apacheHttpClientBuilder.build()
7361
}.toEither
74-
s3Client <- credentialsProvider(config).leftMap(new IllegalArgumentException(_)).flatMap { credsProv =>
75-
Try(
76-
S3Client
77-
.builder()
78-
.overrideConfiguration(overrideConfig)
79-
.serviceConfiguration(s3Config)
80-
.credentialsProvider(credsProv)
81-
.httpClient(httpClient),
82-
).toEither
62+
s3Client <- credentialsProvider(config)
63+
.leftMap(new IllegalArgumentException(_))
64+
.flatMap {
65+
credsProv =>
66+
Try(
67+
S3Client
68+
.builder()
69+
.overrideConfiguration((clientOverrideConfigurationBuilder: ClientOverrideConfiguration.Builder) =>
70+
customiseOverrideConfiguration(config, clientOverrideConfigurationBuilder),
71+
)
72+
.serviceConfiguration(s3Config)
73+
.credentialsProvider(credsProv)
74+
.httpClient(httpClient),
75+
).toEither
8376

84-
}.map { builder =>
85-
config
86-
.region
87-
.fold(builder)(reg => builder.region(Region.of(reg)))
88-
}.map { builder =>
89-
config.customEndpoint.fold(builder)(cE => builder.endpointOverride(URI.create(cE)))
90-
}.flatMap { builder =>
91-
Try(builder.build()).toEither
92-
}
77+
}.map { builder =>
78+
config
79+
.region
80+
.fold(builder)(reg => builder.region(Region.of(reg)))
81+
}.map { builder =>
82+
config.customEndpoint.fold(builder)(cE => builder.endpointOverride(URI.create(cE)))
83+
}.flatMap { builder =>
84+
Try(builder.build()).toEither
85+
}
9386
} yield s3Client
9487

88+
private def customiseOverrideConfiguration(
89+
config: S3ConnectionConfig,
90+
clientOverrideConfigurationBuilder: ClientOverrideConfiguration.Builder,
91+
): Unit = {
92+
clientOverrideConfigurationBuilder
93+
.retryStrategy { (retryStrategyBuilder: RetryStrategy.Builder[_, _]) =>
94+
customiseRetryStrategy(config, retryStrategyBuilder)
95+
}
96+
()
97+
}
98+
99+
private def customiseRetryStrategy(
100+
config: S3ConnectionConfig,
101+
retryStrategyBuilder: RetryStrategy.Builder[_, _],
102+
): Unit = {
103+
retryStrategyBuilder
104+
.maxAttempts(config.httpRetryConfig.getRetryLimit)
105+
106+
retryStrategyBuilder.backoffStrategy(
107+
new FixedDelayWithoutJitter(
108+
Duration.ofMillis(
109+
config.httpRetryConfig.getRetryIntervalMillis,
110+
),
111+
),
112+
)
113+
114+
()
115+
}
116+
95117
private def credentialsFromConfig(awsConfig: S3ConnectionConfig): Either[String, AwsCredentialsProvider] =
96118
awsConfig.accessKey.zip(awsConfig.secretKey) match {
97119
case Some((access, secret)) =>

project/Dependencies.scala

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,7 @@ object Dependencies {
7878
val jerseyCommonVersion = "3.1.8"
7979

8080
val calciteVersion = "1.34.0"
81-
val awsSdkVersion = "2.25.70"
81+
val awsSdkVersion = "2.26.29"
8282

8383
val azureDataLakeVersion = "12.20.0"
8484
val azureIdentityVersion = "1.13.2"
@@ -95,7 +95,7 @@ object Dependencies {
9595
val openCsvVersion = "5.9"
9696
val jsonSmartVersion = "2.5.1"
9797

98-
val xzVersion = "1.9"
98+
val xzVersion = "1.10"
9999
val lz4Version = "1.8.0"
100100

101101
val bouncyCastleVersion = "1.78.1"

0 commit comments

Comments
 (0)