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
相关产品推荐
相关产品推荐

