From ebf673d4c13d903eaccbe4995a2308733b7ddc32 Mon Sep 17 00:00:00 2001 From: Andreas Maechler Date: Thu, 1 Oct 2026 13:44:07 -0600 Subject: [PATCH] Support IAM auth for Kafka supervisors and lookups Kafka supervisors and Kafka lookups could not authenticate to an IAM-enabled Amazon MSK cluster, because neither extension shipped the aws-msk-iam-auth plugin that registers the AWS_MSK_IAM SASL mechanism. Add aws-msk-iam-auth 2.3.9 to both extensions at runtime scope and exclude all of its transitives, so only the plugin jar ships. The AWS SDK v2 it needs is already in lib/ through druid-server and druid-aws-common, and extensions resolve it parent-first, so bundling another copy would put two SDK versions in one JVM. The plugin cannot live in core itself: its callback handlers implement Kafka client interfaces, and kafka-clients is bundled per extension. No code change is needed in either extension. Both pass consumer properties through unfiltered and set the thread context classloader to the extension classloader around consumer construction, which is where Kafka loads the login module and callback handler. The docs list the IAM actions a consumer needs and recommend awsAddDefaultProviders="false" with assume-role, so that a failed assumption does not fall back to the ambient identity. The tests build a consumer through each extension's own factory with the documented properties. Login happens in the constructor, so no broker is needed, and a missing plugin or SDK artifact fails the test. --- docs/ingestion/kafka-ingestion.md | 26 +++++++ docs/querying/kafka-extraction-namespace.md | 4 + .../kafka-extraction-namespace/pom.xml | 12 +++ .../druid/query/lookup/AwsMskIamAuthTest.java | 72 ++++++++++++++++++ .../kafka-indexing-service/pom.xml | 12 +++ .../indexing/kafka/AwsMskIamAuthTest.java | 73 +++++++++++++++++++ licenses.yaml | 24 ++++++ pom.xml | 6 ++ 8 files changed, 229 insertions(+) create mode 100644 extensions-core/kafka-extraction-namespace/src/test/java/org/apache/druid/query/lookup/AwsMskIamAuthTest.java create mode 100644 extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/AwsMskIamAuthTest.java diff --git a/docs/ingestion/kafka-ingestion.md b/docs/ingestion/kafka-ingestion.md index 0441dde69d16..11c6ecc634c7 100644 --- a/docs/ingestion/kafka-ingestion.md +++ b/docs/ingestion/kafka-ingestion.md @@ -196,6 +196,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 9f69c51706d5..000c88b523f7 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 d378875dab1f..d4dc4749a247 100644 --- a/licenses.yaml +++ b/licenses.yaml @@ -3210,6 +3210,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 @@ -4372,6 +4384,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 46cd05e6cd06..79f29f2a8c88 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.1 + 2.3.9 3.2.4 0.8.15 2.0.5 @@ -420,6 +421,11 @@ aws-crt-client ${aws.sdk.v2.version} + + software.amazon.msk + aws-msk-iam-auth + ${aws.msk.iam.auth.version} + net.minidev json-smart