From 9f7c602321d5abfc01a8407e51d9397d11ff664b Mon Sep 17 00:00:00 2001 From: Jared Yu Date: Sun, 26 Jul 2026 23:55:15 -0700 Subject: [PATCH] [filesystems] Fix S3 assumed role ARN not used for file I/O When s3.assumed.role.arn is configured in server.yaml without static credentials, configure AssumedRoleCredentialProvider so that the assumed role is actually used for S3 operations (remote log, KV snapshots, lake offsets). Previously, the code only logged a message and returned without setting the credential provider, causing S3A to use ambient credentials from the default chain instead of the configured role. Closes #3761 --- .../fluss/fs/s3/S3FileSystemPlugin.java | 21 +++++++++++++---- .../fluss/fs/s3/S3FileSystemPluginTest.java | 23 +++++++++++++++++++ 2 files changed, 40 insertions(+), 4 deletions(-) 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();