diff --git a/docs/ingestion/kafka-ingestion.md b/docs/ingestion/kafka-ingestion.md
index 91f1706f134c..264c2d7be982 100644
--- a/docs/ingestion/kafka-ingestion.md
+++ b/docs/ingestion/kafka-ingestion.md
@@ -197,6 +197,32 @@ When you define the consumer properties in the supervisor spec, use the dynamic
When connecting to Kafka, Druid replaces the environment variables with their corresponding values.
+##### Amazon MSK with IAM authentication
+
+The `druid-kafka-indexing-service` extension includes the [Amazon MSK Library for AWS Identity and Access Management](https://github.com/aws/aws-msk-iam-auth), so supervisors can read from an MSK cluster that uses [IAM access control](https://docs.aws.amazon.com/msk/latest/developerguide/iam-access-control.html) with nothing extra to install. Set these consumer properties, pointing `bootstrap.servers` at the cluster's [IAM bootstrap brokers](https://docs.aws.amazon.com/msk/latest/developerguide/msk-get-bootstrap-brokers.html):
+
+```json
+"consumerProperties": {
+ "bootstrap.servers": "b-1.example.c1.kafka.us-east-1.amazonaws.com:9098",
+ "security.protocol": "SASL_SSL",
+ "sasl.mechanism": "AWS_MSK_IAM",
+ "sasl.jaas.config": "software.amazon.msk.auth.iam.IAMLoginModule required;",
+ "sasl.client.callback.handler.class": "software.amazon.msk.auth.iam.IAMClientCallbackHandler"
+}
+```
+
+Grant the IAM identity that Druid uses the actions AWS lists for [consuming data](https://docs.aws.amazon.com/msk/latest/developerguide/iam-access-control-use-cases.html): `kafka-cluster:Connect`, `kafka-cluster:DescribeTopic`, `kafka-cluster:ReadData`, `kafka-cluster:DescribeGroup`, and `kafka-cluster:AlterGroup`. [Scope each action](https://docs.aws.amazon.com/msk/latest/developerguide/kafka-actions.html) to the cluster, topic, or group it applies to.
+
+Druid finds credentials through the default AWS credentials provider chain. The Overlord and the Peons or Indexers that run Kafka tasks pick up an EC2 instance profile, a Kubernetes service account role, or environment credentials without extra configuration.
+
+To assume a role, add `awsRoleArn` and `awsStsRegion` to `sasl.jaas.config`. Also add `awsAddDefaultProviders="false"`. Without it, a failed role assumption falls back to the default credentials, and Druid connects as that identity instead:
+
+```json
+"sasl.jaas.config": "software.amazon.msk.auth.iam.IAMLoginModule required awsRoleArn=\"arn:aws:iam::123456789012:role/msk-consumer\" awsStsRegion=\"us-east-1\" awsAddDefaultProviders=\"false\";"
+```
+
+Kafka lookups take the same properties in `kafkaProperties`. See [Kafka lookups](../querying/kafka-extraction-namespace.md#amazon-msk-with-iam-authentication).
+
#### Idle configuration
:::info
diff --git a/docs/querying/kafka-extraction-namespace.md b/docs/querying/kafka-extraction-namespace.md
index f5c91e7847d7..c2abd040a966 100644
--- a/docs/querying/kafka-extraction-namespace.md
+++ b/docs/querying/kafka-extraction-namespace.md
@@ -75,6 +75,10 @@ This input topic would be consumed from the beginning, and result in a lookup na
Now when a query uses this extraction namespace, the country codes can be mapped to the full country name at query time.
+## Amazon MSK with IAM authentication
+
+The `druid-kafka-extraction-namespace` extension includes the [Amazon MSK Library for AWS Identity and Access Management](https://github.com/aws/aws-msk-iam-auth), so a lookup can read from an MSK cluster that uses IAM access control. Put the properties from [Kafka ingestion](../ingestion/kafka-ingestion.md#amazon-msk-with-iam-authentication) in `kafkaProperties`. Every service that loads the lookup needs IAM credentials with the permissions listed there.
+
## Tombstones and Deleting Records
The Kafka lookup extractor treats `null` Kafka messages as tombstones. This means that a record on the input topic with a `null` message payload on Kafka will remove the associated key from the lookup map, effectively deleting it.
diff --git a/extensions-core/kafka-extraction-namespace/pom.xml b/extensions-core/kafka-extraction-namespace/pom.xml
index 44584a5da109..cf5705aaeb4a 100644
--- a/extensions-core/kafka-extraction-namespace/pom.xml
+++ b/extensions-core/kafka-extraction-namespace/pom.xml
@@ -67,6 +67,18 @@
+
+
+ software.amazon.msk
+ aws-msk-iam-auth
+ runtime
+
+
+ *
+ *
+
+
+
com.google.code.findbugs
jsr305
diff --git a/extensions-core/kafka-extraction-namespace/src/test/java/org/apache/druid/query/lookup/AwsMskIamAuthTest.java b/extensions-core/kafka-extraction-namespace/src/test/java/org/apache/druid/query/lookup/AwsMskIamAuthTest.java
new file mode 100644
index 000000000000..8bf9fe4d3b71
--- /dev/null
+++ b/extensions-core/kafka-extraction-namespace/src/test/java/org/apache/druid/query/lookup/AwsMskIamAuthTest.java
@@ -0,0 +1,72 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.druid.query.lookup;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import javax.security.sasl.Sasl;
+import java.util.Map;
+
+/**
+ * Checks that a lookup configured for {@code AWS_MSK_IAM} as documented logs in using only the AWS SDK on Druid's
+ * core classpath. Login happens in the consumer constructor, so no broker is needed.
+ */
+public class AwsMskIamAuthTest
+{
+ private static final String LOGIN_MODULE = "software.amazon.msk.auth.iam.IAMLoginModule";
+
+ @Test
+ public void testDefaultCredentials() throws Exception
+ {
+ createAndCloseConsumer(LOGIN_MODULE + " required;");
+
+ // Construction never connects, so separately check the mechanism the SASL handshake would look up.
+ Assertions.assertNotNull(
+ Sasl.createSaslClient(new String[]{"AWS_MSK_IAM"}, null, "kafka", "localhost", Map.of(), callbacks -> {}),
+ "AWS_MSK_IAM SASL mechanism is not registered"
+ );
+ }
+
+ /**
+ * Assume-role builds an {@code StsClient} at login, which reaches more of the SDK than the default chain.
+ */
+ @Test
+ public void testAssumeRole()
+ {
+ createAndCloseConsumer(
+ LOGIN_MODULE + " required awsRoleArn=\"arn:aws:iam::123456789012:role/msk-consumer\" awsStsRegion=\"us-east-1\""
+ + " awsAddDefaultProviders=\"false\";"
+ );
+ }
+
+ private static void createAndCloseConsumer(final String jaasConfig)
+ {
+ final Map kafkaProperties = Map.of(
+ "bootstrap.servers", "localhost:9098",
+ "security.protocol", "SASL_SSL",
+ "sasl.mechanism", "AWS_MSK_IAM",
+ "sasl.jaas.config", jaasConfig,
+ "sasl.client.callback.handler.class", "software.amazon.msk.auth.iam.IAMClientCallbackHandler"
+ );
+
+ new KafkaLookupExtractorFactory(null, "lookup-topic", kafkaProperties).getConsumer().close();
+ }
+}
diff --git a/extensions-core/kafka-indexing-service/pom.xml b/extensions-core/kafka-indexing-service/pom.xml
index cb508bc8d850..d37bee7444cf 100644
--- a/extensions-core/kafka-indexing-service/pom.xml
+++ b/extensions-core/kafka-indexing-service/pom.xml
@@ -67,6 +67,18 @@
+
+
+ software.amazon.msk
+ aws-msk-iam-auth
+ runtime
+
+
+ *
+ *
+
+
+
com.google.code.findbugs
jsr305
diff --git a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/AwsMskIamAuthTest.java b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/AwsMskIamAuthTest.java
new file mode 100644
index 000000000000..1b1d5ce02f13
--- /dev/null
+++ b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/AwsMskIamAuthTest.java
@@ -0,0 +1,73 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.druid.indexing.kafka;
+
+import org.apache.druid.jackson.DefaultObjectMapper;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import javax.security.sasl.Sasl;
+import java.util.Map;
+
+/**
+ * Checks that a consumer configured for {@code AWS_MSK_IAM} as documented logs in using only the AWS SDK on Druid's
+ * core classpath. Login happens in the consumer constructor, so no broker is needed.
+ */
+public class AwsMskIamAuthTest
+{
+ private static final String LOGIN_MODULE = "software.amazon.msk.auth.iam.IAMLoginModule";
+
+ @Test
+ public void testDefaultCredentials() throws Exception
+ {
+ createAndCloseConsumer(LOGIN_MODULE + " required;");
+
+ // Construction never connects, so separately check the mechanism the SASL handshake would look up.
+ Assertions.assertNotNull(
+ Sasl.createSaslClient(new String[]{"AWS_MSK_IAM"}, null, "kafka", "localhost", Map.of(), callbacks -> {}),
+ "AWS_MSK_IAM SASL mechanism is not registered"
+ );
+ }
+
+ /**
+ * Assume-role builds an {@code StsClient} at login, which reaches more of the SDK than the default chain.
+ */
+ @Test
+ public void testAssumeRole()
+ {
+ createAndCloseConsumer(
+ LOGIN_MODULE + " required awsRoleArn=\"arn:aws:iam::123456789012:role/msk-consumer\" awsStsRegion=\"us-east-1\""
+ + " awsAddDefaultProviders=\"false\";"
+ );
+ }
+
+ private static void createAndCloseConsumer(final String jaasConfig)
+ {
+ final Map consumerProperties = Map.of(
+ "bootstrap.servers", "localhost:9098",
+ "security.protocol", "SASL_SSL",
+ "sasl.mechanism", "AWS_MSK_IAM",
+ "sasl.jaas.config", jaasConfig,
+ "sasl.client.callback.handler.class", "software.amazon.msk.auth.iam.IAMClientCallbackHandler"
+ );
+
+ KafkaRecordSupplier.getKafkaConsumer(new DefaultObjectMapper(), consumerProperties, null).close();
+ }
+}
diff --git a/licenses.yaml b/licenses.yaml
index 9f0d1ad41721..3e865395c958 100644
--- a/licenses.yaml
+++ b/licenses.yaml
@@ -3241,6 +3241,18 @@ notices:
---
+name: Amazon MSK Library for AWS Identity and Access Management
+license_category: binary
+module: extensions/druid-kafka-indexing-service
+license_name: Apache License version 2.0
+version: 2.3.9
+libraries:
+ - software.amazon.msk: aws-msk-iam-auth
+notice: |
+ Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
+
+---
+
name: Apache Parquet
license_category: binary
@@ -4403,6 +4415,18 @@ notices:
---
+name: Amazon MSK Library for AWS Identity and Access Management
+license_category: binary
+module: extensions/druid-kafka-extraction-namespace
+license_name: Apache License version 2.0
+version: 2.3.9
+libraries:
+ - software.amazon.msk: aws-msk-iam-auth
+notice: |
+ Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
+
+---
+
name: Microsoft Azure SDK For Key Vault Core
license_category: binary
module: extensions/druid-azure-extensions
diff --git a/pom.xml b/pom.xml
index efd4ea828ae0..b2b437550e3d 100644
--- a/pom.xml
+++ b/pom.xml
@@ -128,6 +128,7 @@
It should be removed once those extensions are upgraded or dropped (see #19109). -->
1.12.797
2.55.4
+ 2.3.9
3.3.0
0.8.15
2.0.5
@@ -416,6 +417,11 @@
aws-crt-client
${aws.sdk.v2.version}
+
+ software.amazon.msk
+ aws-msk-iam-auth
+ ${aws.msk.iam.auth.version}
+
net.minidev
json-smart