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

NodeJS通过带SASL扩展的OpenID连接Confluent Cloud Kafka遇认证问题

问题:NodeJS连接Confluent Cloud Kafka时SASL扩展参数无效

背景

我正在用NodeJS连接配置了OpenID身份提供商的Confluent Cloud Kafka服务器,使用kafka-console-consumer命令可以正常连接,但基于KafkaJS实现时,添加SASL扩展参数(logicalCluster和identityPoolId)遇到了认证失败的问题,报错提示logicalCluster: CLUSTER_ID_MISSING_OR_EMPTY。

可正常工作的命令行配置

Kafka配置文件(kafka.properties)

# ./kafka.properties
security.protocol=SASL_SSL
sasl.mechanism=OAUTHBEARER
sasl.login.callback.handler.class=org.apache.kafka.common.security.oauthbearer.secured.OAuthBearerLoginCallbackHandler
sasl.login.connect.timeout.ms=15000
sasl.oauthbearer.token.endpoint.url=https://myidp.example.com/oauth2/default/v1/token
sasl.jaas.config=org.apache.kafka.common.security.oauthbearer.OAuthBearerLoginModule required \
    clientId="MY_CLIENT_ID" \
    clientSecret="MY_CLIENT_SECRET" \
    scope="SCOPE" \
    extension_logicalCluster="CLUSTER_ID" \
    extension_identityPoolId="IDENTITY_POOL_ID" ;

消费者启动命令

kafka-console-consumer \
        --bootstrap-server KAFKA_BROKER \
        --consumer.config ./kafka.properties \
        --topic MY_TOPIC \
        --group GROUP_NAME

当前KafkaJS实现代码

# oauthBearerProviderOpenId.ts
import { Issuer, TokenSet } from 'openid-client';

export interface oauthBearerProviderOpenIdOptions {
    issuer: string;
    clientId: string;
    clientSecret: string;
    refreshThresholdMs: number;
    scope: string,
    logicalCluster: string,
    identityPoolId: string,
}

export const oauthBearerProviderOpenId = async (options: oauthBearerProviderOpenIdOptions) => {
    const issuer = await Issuer.discover(options.issuer);
    const { Client } = issuer;

    const client = new Client({
        client_id: options.clientId,
        client_secret: options.clientSecret,
    });

    let tokenPromise: Promise<string>;
    let tokenSet: TokenSet | null;

    async function refreshToken() {
        try {
            if (tokenSet == null) {
                tokenSet = await client.grant({
                    grant_type: 'client_credentials',
                    scope: options.scope,
                });
            }

            setTimeout(() => {
                tokenPromise = refreshToken()
            }, tokenSet!.expires_in);

            if(!tokenSet.access_token) {
                throw new Error('Unable to fetch access_token');
            }
            
            return tokenSet.access_token;

        } catch (error) {
            const e = error as any;
            tokenSet = null;
            console.error(e.data.payload.toString());
            throw error;
        }
    }

    tokenPromise = refreshToken();

    return async function () {
        return {
            value: await tokenPromise
        }
    }
};

错误日志

SASL OAUTHBEARER authentication failed: Authentication failed: 1 extensions are invalid! They are: logicalCluster: CLUSTER_ID_MISSING_OR_EMPTY
KafkaJSSASLAuthenticationError: SASL OAUTHBEARER authentication failed: Authentication failed: 1 extensions are invalid! They are: logicalCluster: CLUSTER_ID_MISSING_OR_EMPTY

解决方法

问题核心是KafkaJS的OAuthBearer提供者返回的对象缺少了SASL扩展参数。Confluent Cloud要求在SASL OAUTHBEARER握手阶段传递logicalCluster和identityPoolId这两个扩展字段,你需要修改返回的token对象,添加extensions属性:

修改token返回部分

将原代码中返回token的函数修改为:

return async function () {
    return {
        value: await tokenPromise,
        extensions: {
            logicalCluster: options.logicalCluster,
            identityPoolId: options.identityPoolId
        }
    }
}

优化token刷新时机(可选)

原代码中直接使用tokenSet.expires_in作为刷新延迟,可能会导致token过期后才触发刷新。建议结合传入的refreshThresholdMs提前刷新:

setTimeout(() => {
    tokenPromise = refreshToken()
}, (tokenSet!.expires_in * 1000) - options.refreshThresholdMs);

这样修改后,KafkaJS客户端会在SASL握手时将扩展参数发送给Confluent Cloud集群,解决认证失败的问题。

内容的提问来源于stack exchange,提问作者Daniel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 01:42:21