diff --git a/fluss-filesystems/fluss-fs-s3/src/main/java/org/apache/fluss/fs/s3/S3FileSystemPlugin.java b/fluss-filesystems/fluss-fs-s3/src/main/java/org/apache/fluss/fs/s3/S3FileSystemPlugin.java index 7e293f5bd6..d9e261f692 100644 --- a/fluss-filesystems/fluss-fs-s3/src/main/java/org/apache/fluss/fs/s3/S3FileSystemPlugin.java +++ b/fluss-filesystems/fluss-fs-s3/src/main/java/org/apache/fluss/fs/s3/S3FileSystemPlugin.java @@ -56,6 +56,13 @@ public class S3FileSystemPlugin implements FileSystemPlugin { private static final String ROLE_ARN_KEY = "fs.s3a.assumed.role.arn"; + /** + * The Hadoop S3A AssumedRoleCredentialProvider class that uses the configured role ARN to + * assume an IAM role for S3 access. + */ + private static final String ASSUMED_ROLE_CREDENTIAL_PROVIDER = + "org.apache.hadoop.fs.s3a.auth.AssumedRoleCredentialProvider"; + private static final String[][] MIRRORED_CONFIG_KEYS = { {"fs.s3a.access-key", "fs.s3a.access.key"}, {"fs.s3a.secret-key", "fs.s3a.secret.key"}, @@ -163,10 +170,16 @@ private void setCredentialProvider(org.apache.hadoop.conf.Configuration hadoopCo } if (hasStaticKeys || hasRoleArn) { - LOG.info( - hasStaticKeys - ? "Using provided static credentials." - : "Using default AWS credential chain with AssumeRole."); + if (hasRoleArn && !hasStaticKeys) { + // When only role ARN is configured, use the AssumedRoleCredentialProvider + // to assume the role for all S3 operations (reads/writes). + hadoopConfig.set(PROVIDER_CONFIG_NAME, ASSUMED_ROLE_CREDENTIAL_PROVIDER); + LOG.info( + "Using AssumedRoleCredentialProvider with role ARN for S3 access: {}", + hadoopConfig.get(ROLE_ARN_KEY)); + } else { + LOG.info("Using provided static credentials."); + } return; } diff --git a/fluss-filesystems/fluss-fs-s3/src/test/java/org/apache/fluss/fs/s3/S3FileSystemPluginTest.java b/fluss-filesystems/fluss-fs-s3/src/test/java/org/apache/fluss/fs/s3/S3FileSystemPluginTest.java index d946348551..4584bdeaa2 100644 --- a/fluss-filesystems/fluss-fs-s3/src/test/java/org/apache/fluss/fs/s3/S3FileSystemPluginTest.java +++ b/fluss-filesystems/fluss-fs-s3/src/test/java/org/apache/fluss/fs/s3/S3FileSystemPluginTest.java @@ -62,10 +62,33 @@ void testServerModeWithRoleArnOnly() { org.apache.hadoop.conf.Configuration hadoopConfig = plugin.buildHadoopConfiguration(flussConfig); + // When only role ARN is configured, AssumedRoleCredentialProvider should be used String providers = hadoopConfig.get(PROVIDER_CONFIG, ""); + assertThat(providers) + .isEqualTo("org.apache.hadoop.fs.s3a.auth.AssumedRoleCredentialProvider"); assertThat(providers).doesNotContain(DynamicTemporaryAWSCredentialsProvider.NAME); } + @Test + void testServerModeWithStaticKeysAndRoleArn() { + // When both static keys and role ARN are provided, static keys should take precedence + Configuration flussConfig = new Configuration(); + flussConfig.setString("fs.s3a.access.key", "testAccessKey"); + flussConfig.setString("fs.s3a.secret.key", "testSecretKey"); + flussConfig.setString( + "fs.s3a.assumed.role.arn", "arn:aws:iam::123456789012:role/test-role"); + + S3FileSystemPlugin plugin = new S3FileSystemPlugin(); + org.apache.hadoop.conf.Configuration hadoopConfig = + plugin.buildHadoopConfiguration(flussConfig); + + // Static keys take precedence, AssumedRoleCredentialProvider should NOT be set + String providers = hadoopConfig.get(PROVIDER_CONFIG, ""); + assertThat(providers).doesNotContain(DynamicTemporaryAWSCredentialsProvider.NAME); + assertThat(providers) + .doesNotContain("org.apache.hadoop.fs.s3a.auth.AssumedRoleCredentialProvider"); + } + @Test void testServerModeWithConfiguredCredentialProvider() { Configuration flussConfig = new Configuration();