Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 26 additions & 0 deletions docs/ingestion/kafka-ingestion.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 4 additions & 0 deletions docs/querying/kafka-extraction-namespace.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
12 changes: 12 additions & 0 deletions extensions-core/kafka-extraction-namespace/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,18 @@
</exclusion>
</exclusions>
</dependency>
<!-- Provides the AWS_MSK_IAM SASL mechanism; its AWS SDK comes from druid-aws-common, so no transitives ship. -->
<dependency>
<groupId>software.amazon.msk</groupId>
<artifactId>aws-msk-iam-auth</artifactId>
<scope>runtime</scope>
<exclusions>
<exclusion>
<groupId>*</groupId>
<artifactId>*</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>com.google.code.findbugs</groupId>
<artifactId>jsr305</artifactId>
Expand Down
Original file line number Diff line number Diff line change
@@ -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<String, String> 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();
}
}
12 changes: 12 additions & 0 deletions extensions-core/kafka-indexing-service/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,18 @@
</exclusion>
</exclusions>
</dependency>
<!-- Provides the AWS_MSK_IAM SASL mechanism; its AWS SDK comes from druid-aws-common, so no transitives ship. -->
<dependency>
<groupId>software.amazon.msk</groupId>
<artifactId>aws-msk-iam-auth</artifactId>
<scope>runtime</scope>
<exclusions>
<exclusion>
<groupId>*</groupId>

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Default-provider SSO profiles fail at runtime

Finding: Both extension POMs exclude every transitive of aws-msk-iam-auth 2.3.9, but this release's ProfileCredentialsProvider loads software.amazon.awssdk.services.sso.auth.SsoProfileCredentialsProviderFactory for profiles using SSO fields and declares the sso and ssooidc modules explicitly. druid-aws-common in this tree supplies the shared auth, regions, and STS modules but not those SSO modules, so a user relying on the default AWS provider chain with an SSO-backed profile cannot resolve credentials and the Kafka consumer fails before authenticating.

Suggestion: Add compatible AWS SDK sso and ssooidc runtime modules to the shared classpath, or retain these plugin transitives in both extension packaging paths, and cover an SSO-backed profile in a smoke test.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks, this is deliberate. The gap is not specific to this PR. Druid's own chain (AWSCredentialsUtils.defaultAWSCredentialsProviderChain, used by S3 and Kinesis) fails on the same sso-session profile with the same "the 'sso' service module must be on the class path" error, because lib/ lacks these modules.

I had Claude look up some specifics around how AWS SSO works:

SDK SSO profiles (sso_session, sso_start_url) use a token from a user's IAM Identity Center sign-in, and ssooidc only refreshes that user token. Workloads use other sources: web identity federation (OIDC token via STS AssumeRoleWithWebIdentity, as with IRSA), credential_process (for example IAM Roles Anywhere), instance profiles, pod identity and task roles. I verified that web identity and credential_process resolve with the shipped lib/. No issue has asked for SSO support, so I left these modules out.

Keeping the plugin's transitives in the Kafka extensions would not fix it. ProfileCredentialsUtils in lib/ loads SsoProfileCredentialsProviderFactory by name from the thread context classloader or its own. Credentials resolve during the SASL handshake on the polling thread, outside the context-classloader switch around consumer construction, so an extension copy is never found. In a reproduction with an extension classloader over the shipped lib/, it worked only when the context classloader was the extension's.

That said, users who need SSO can add sso, ssooidc and aws-json-protocol to lib/ in their own Druid install. sso needs aws-json-protocol; without it, resolution fails with NoClassDefFoundError. I verified this layout for both the plugin and Druid's chain. We could always add them to druid-aws-common if there is demand (in a separate PR).

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Follow-up assessment

Thanks for the clarification. I rechecked the current head across all 8 changed files, including both extension POMs and Druid's existing AWSCredentialsUtils.defaultAWSCredentialsProviderChain. The missing SSO modules are already absent from the shared druid-aws-common classpath, and this PR does not change that path, so the SSO limitation is not a PR-caused regression. I am withdrawing the P2 finding; no further action is needed on this thread.

Reviewed 8 of 8 changed files; this follow-up has 0 remaining findings. Static git diff --check e855cd9c3316710361c7f2cf0b1b8331141c20ad ebf673d4c13d903eaccbe4995a2308733b7ddc32 passed; no tests or builds were run.

<artifactId>*</artifactId>
</exclusion>
</exclusions>
</dependency>
<dependency>
<groupId>com.google.code.findbugs</groupId>
<artifactId>jsr305</artifactId>
Expand Down
Original file line number Diff line number Diff line change
@@ -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<String, Object> 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();
}
}
24 changes: 24 additions & 0 deletions licenses.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,7 @@
It should be removed once those extensions are upgraded or dropped (see #19109). -->
<aws.sdk.v1.version>1.12.797</aws.sdk.v1.version>
<aws.sdk.v2.version>2.55.4</aws.sdk.v2.version>
<aws.msk.iam.auth.version>2.3.9</aws.msk.iam.auth.version>
<caffeine.version>3.3.0</caffeine.version>
<jacoco.version>0.8.15</jacoco.version>
<testcontainers.version>2.0.5</testcontainers.version>
Expand Down Expand Up @@ -416,6 +417,11 @@
<artifactId>aws-crt-client</artifactId>
<version>${aws.sdk.v2.version}</version>
</dependency>
<dependency>
<groupId>software.amazon.msk</groupId>
<artifactId>aws-msk-iam-auth</artifactId>
<version>${aws.msk.iam.auth.version}</version>
</dependency>
<dependency>
<groupId>net.minidev</groupId>
<artifactId>json-smart</artifactId>
Expand Down
Loading