You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何为Apache Kafka集群配置Azure AD OAuth2认证及Java验证Azure访问令牌

Hey there, let's tackle your two related problems: setting up Azure AD-based OAuth2 authentication for your Apache Kafka cluster, and validating Azure Access Tokens in Java since Azure AD doesn't provide an introspect endpoint. Here's a step-by-step breakdown:

1. Configure Apache Kafka Cluster for Azure AD OAuth2 Authentication

Kafka's OAuth2 support can be adapted to work with Azure AD by leveraging JWT validation (since Azure AD's access tokens are standard JWTs) instead of the introspect endpoint. Here's how to set it up:

1.1 Broker Configuration (server.properties)

First, enable SASL OAuth2 as the authentication mechanism and configure the necessary JAAS settings:

# Enable SASL listener
listeners=PLAINTEXT://:9092,SASL_PLAINTEXT://:9093
advertised.listeners=PLAINTEXT://localhost:9092,SASL_PLAINTEXT://localhost:9093
security.inter.broker.protocol=SASL_PLAINTEXT

# Enable OAuthBEARER mechanism
sasl.enabled.mechanisms=OAUTHBEARER
sasl.mechanism.inter.broker.protocol=OAUTHBEARER

# JAAS config for broker's OAuth2 login
listener.name.sasl_plaintext.oauthbearer.sasl.jaas.config=org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required \
  oauth.client.id="your-azure-ad-client-id" \
  oauth.client.secret="your-azure-ad-client-secret" \
  oauth.token.endpoint.uri="https://login.microsoftonline.com/your-tenant-id/oauth2/v2.0/token";

# Custom callback handler for Azure AD token validation (we'll build this next)
listener.name.sasl_plaintext.oauthbearer.sasl.server.callback.handler.class=com.yourorg.kafka.AzureAdOauthServerCallbackHandler

1.2 Client Configuration (Producer/Consumer Properties)

For Kafka producers and consumers, configure them to use OAuth2 to authenticate with the broker:

security.protocol=SASL_PLAINTEXT
sasl.mechanism=OAUTHBEARER
sasl.jaas.config=org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required \
  oauth.client.id="your-azure-ad-client-id" \
  oauth.client.secret="your-azure-ad-client-secret" \
  oauth.token.endpoint.uri="https://login.microsoftonline.com/your-tenant-id/oauth2/v2.0/token";
sasl.login.callback.handler.class=org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginCallbackHandler
2. Java Solution to Validate Azure Access Tokens

Since Azure AD lacks an introspect endpoint, we'll validate the JWT token directly by checking its signature, issuer, audience, and expiration. Here are two reliable approaches:

2.1 Use Microsoft's Official MSAL4J Library

MSAL4J is the recommended way to handle Azure AD token operations in Java—it automatically fetches public keys and validates all standard JWT claims.

Step 1: Add Dependency (Maven)

<dependency>
    <groupId>com.microsoft.azure</groupId>
    <artifactId>msal4j</artifactId>
    <version>1.23.0</version>
</dependency>

Step 2: Token Validation Code

import com.microsoft.aad.msal4j.Jwt;
import com.microsoft.aad.msal4j.ValidationParameters;
import java.util.Collections;
import java.util.concurrent.CompletableFuture;

public class AzureAdTokenValidator {
    private static final String TENANT_ID = "your-tenant-id";
    private static final String KAFKA_CLIENT_ID = "your-kafka-client-id";
    private static final String ISSUER = "https://login.microsoftonline.com/" + TENANT_ID + "/v2.0";

    public static boolean isTokenValid(String accessToken) {
        try {
            ValidationParameters validationParams = ValidationParameters.builder()
                    .issuer(ISSUER)
                    .audience(Collections.singletonList(KAFKA_CLIENT_ID))
                    .build();

            CompletableFuture<Boolean> isValid = Jwt.decode(accessToken)
                    .validate(validationParams);

            return isValid.get();
        } catch (Exception e) {
            System.err.println("Token validation failed: " + e.getMessage());
            return false;
        }
    }
}

2.2 Manual Validation with JJWT Library

If you prefer a lightweight approach without MSAL4J, use the JJWT library to parse and validate the JWT:

Step 1: Add Dependencies (Maven)

<dependency>
    <groupId>io.jsonwebtoken</groupId>
    <artifactId>jjwt-api</artifactId>
    <version>0.11.5</version>
</dependency>
<dependency>
    <groupId>io.jsonwebtoken</groupId>
    <artifactId>jjwt-impl</artifactId>
    <version>0.11.5</version>
    <scope>runtime</scope>
</dependency>
<dependency>
    <groupId>io.jsonwebtoken</groupId>
    <artifactId>jjwt-jackson</artifactId>
    <version>0.11.5</version>
    <scope>runtime</scope>
</dependency>

Step 2: Token Validation Code

import io.jsonwebtoken.Claims;
import io.jsonwebtoken.Jwts;
import io.jsonwebtoken.security.Keys;
import java.net.URL;
import java.security.PublicKey;

public class JwtTokenValidator {
    private static final String TENANT_ID = "your-tenant-id";
    private static final String KAFKA_CLIENT_ID = "your-kafka-client-id";
    private static final String ISSUER = "https://login.microsoftonline.com/" + TENANT_ID + "/v2.0";
    private static final String JWKS_URI = "https://login.microsoftonline.com/" + TENANT_ID + "/discovery/v2.0/keys";

    public static boolean isTokenValid(String accessToken) {
        try {
            // Fetch Azure AD public keys from JWKS endpoint
            PublicKey publicKey = Keys.publicKeyFromJwkSet(Jwts.SIG.RS256.jwkSetFromUrl(new URL(JWKS_URI)), Jwts.SIG.RS256);

            // Parse and validate token
            Claims claims = Jwts.parser()
                    .verifyWith(publicKey)
                    .requireIssuer(ISSUER)
                    .requireAudience(KAFKA_CLIENT_ID)
                    .build()
                    .parseSignedClaims(accessToken)
                    .getPayload();

            // Add any custom validation logic here (e.g., check roles)
            return true;
        } catch (Exception e) {
            System.err.println("Token validation failed: " + e.getMessage());
            return false;
        }
    }
}

2.3 Integrate Validation into Kafka's Callback Handler

Now, wire up the token validation into the custom callback handler we referenced in the broker config:

import org.apache.kafka.common.security.oauthbearer.OAuthBearerValidatorCallback;
import javax.security.auth.callback.Callback;
import javax.security.auth.callback.CallbackHandler;
import javax.security.auth.callback.UnsupportedCallbackException;
import java.io.IOException;

public class AzureAdOauthServerCallbackHandler implements CallbackHandler {
    @Override
    public void handle(Callback[] callbacks) throws IOException, UnsupportedCallbackException {
        for (Callback callback : callbacks) {
            if (callback instanceof OAuthBearerValidatorCallback) {
                OAuthBearerValidatorCallback validatorCallback = (OAuthBearerValidatorCallback) callback;
                String token = validatorCallback.tokenValue();

                // Use either of our validation methods
                boolean isValid = AzureAdTokenValidator.isTokenValid(token);
                
                if (isValid) {
                    validatorCallback.validated(true);
                    // Optional: Set user/principal from JWT claims
                    // validatorCallback.principal(claims.get("upn"));
                } else {
                    validatorCallback.error("Invalid Azure AD access token");
                }
            } else {
                throw new UnsupportedCallbackException(callback);
            }
        }
    }
}

内容的提问来源于stack exchange,提问作者Vinod K kumar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.06 19:57:42