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

如何用C#/C++/librdkafka结合Azure AD实现Kafka OAuthBearer认证

解决方案:C#/C++ 通过Azure AD OAuth令牌向Apache Kafka生产消息

C# 实现(基于Confluent.Kafka)

配置误区说明

sasl.login.callback.handler.class是Java Kafka客户端特有的配置项,用于指定自定义登录回调处理器。C#使用的Confluent.Kafka客户端基于librdkafka,完全不需要这个配置,只需正确配置OAuth Bearer相关参数即可。

完整代码示例(内置令牌获取)

客户端会自动从Azure AD获取并刷新令牌,无需手动处理:

using Confluent.Kafka;
using System;
using System.Threading.Tasks;

public class KafkaAzureAdProducer
{
    public static async Task ProduceMessageAsync()
    {
        string server = "your-kafka-bootstrap-servers:9092";
        string clientId = "your-client-id";
        string topic = "your-target-topic";

        var config = new ProducerConfig
        {
            BootstrapServers = server,
            ClientId = clientId,
            SecurityProtocol = SecurityProtocol.SaslPlaintext,
            SaslMechanism = SaslMechanism.OAuthBearer,
            // Azure AD OAuth核心配置
            SaslOauthbearerClientId = "MyClientId",
            SaslOauthbearerClientSecret = "MyClientSecret",
            SaslOauthbearerScope = "ServerClientId/.default",
            SaslOauthbearerTokenEndpointUrl = "https://login.microsoftonline.com/MyTenantId/oauth2/v2.0/token"
        };

        using (var producer = new ProducerBuilder<Null, string>(config).Build())
        {
            try
            {
                var result = await producer.ProduceAsync(topic, new Message<Null, string> { Value = "Hello from Azure AD OAuth!" });
                Console.WriteLine($"消息已发送到分区 {result.Partition},偏移量 {result.Offset}");
            }
            catch (ProduceException<Null, string> ex)
            {
                Console.WriteLine($"发送失败: {ex.Error.Reason}");
            }
        }
    }
}

自定义令牌获取(可选)

如果需要使用托管身份、自定义令牌缓存等逻辑,可通过SaslOauthbearerTokenRefreshCallback实现:

var config = new ProducerConfig
{
    BootstrapServers = server,
    ClientId = clientId,
    SecurityProtocol = SecurityProtocol.SaslPlaintext,
    SaslMechanism = SaslMechanism.OAuthBearer,
    // 设置自定义令牌刷新回调
    SaslOauthbearerTokenRefreshCallback = async (tokenData, cancellationToken) =>
    {
        var azureToken = await GetAzureAdTokenAsync();
        tokenData.Token = azureToken.AccessToken;
        // 设置令牌过期时间(毫秒级时间戳)
        tokenData.ExpirationTimestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds() + (long)azureToken.ExpiresIn.TotalMilliseconds;
    }
};

// 自定义Azure AD令牌获取方法
private static async Task<AzureAdToken> GetAzureAdTokenAsync()
{
    // 可使用Azure.Identity库实现,示例为ClientCredential模式
    var client = new Azure.Identity.SecretClient(
        new Uri("https://login.microsoftonline.com/MyTenantId"),
        new Azure.Identity.ClientCredential("MyClientId", "MyClientSecret"));
    var response = await client.GetTokenAsync(new Azure.Core.TokenRequestContext(new[] { "ServerClientId/.default" }));
    
    return new AzureAdToken
    {
        AccessToken = response.Token,
        ExpiresIn = response.ExpiresOn - DateTimeOffset.UtcNow
    };
}

public class AzureAdToken
{
    public string AccessToken { get; set; }
    public TimeSpan ExpiresIn { get; set; }
}

C++ 实现(基于librdkafka)

librdkafka同样不需要sasl.login.callback.handler.class配置,直接通过OAuth相关参数完成Azure AD认证:

完整代码示例

#include <librdkafka/rdkafkacpp.h>
#include <iostream>
#include <string>

// 交付报告回调类
class DeliveryReportCb : public RdKafka::DeliveryReportCb {
public:
    void dr_cb(RdKafka::Message &message) override {
        if (message.err()) {
            std::cerr << "发送失败: " << message.errstr() << std::endl;
        } else {
            std::cout << "消息已发送到分区 " << message.partition() << ", 偏移量 " << message.offset() << std::endl;
        }
    }
};

int main() {
    std::string brokers = "your-kafka-bootstrap-servers:9092";
    std::string topic = "your-target-topic";

    RdKafka::Conf *conf = RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL);
    std::string errstr;

    // 基础配置
    conf->set("bootstrap.servers", brokers, errstr);
    conf->set("client.id", "your-client-id", errstr);
    conf->set("security.protocol", "sasl_plaintext", errstr);
    conf->set("sasl.mechanism", "OAUTHBEARER", errstr);

    // Azure AD OAuth配置
    conf->set("sasl.oauthbearer.client.id", "MyClientId", errstr);
    conf->set("sasl.oauthbearer.client.secret", "MyClientSecret", errstr);
    conf->set("sasl.oauthbearer.scope", "ServerClientId/.default", errstr);
    conf->set("sasl.oauthbearer.token.endpoint.url", "https://login.microsoftonline.com/MyTenantId/oauth2/v2.0/token", errstr);

    // 设置交付报告回调
    DeliveryReportCb dr_cb;
    conf->set("dr_cb", &dr_cb, errstr);

    // 创建生产者实例
    RdKafka::Producer *producer = RdKafka::Producer::create(conf, errstr);
    if (!producer) {
        std::cerr << "创建生产者失败: " << errstr << std::endl;
        delete conf;
        return 1;
    }

    // 发送消息
    std::string message = "Hello from Azure AD OAuth (C++)!";
    RdKafka::ErrorCode err = producer->produce(
        topic, RdKafka::Topic::PARTITION_UA,
        RdKafka::Producer::RK_MSG_COPY,
        const_cast<char *>(message.c_str()), message.size(),
        nullptr, nullptr, nullptr);

    if (err != RdKafka::ERR_NO_ERROR) {
        std::cerr << "生产消息失败: " << RdKafka::err2str(err) << std::endl;
    } else {
        std::cout << "消息已排入队列" << std::endl;
    }

    // 等待所有消息发送完成
    producer->flush(5000);

    delete producer;
    delete conf;
    return 0;
}

自定义令牌获取(可选)

若需自定义令牌逻辑,可通过设置sasl.oauthbearer.token.refresh.callback回调函数实现,具体参考librdkafka官方文档中OAUTHBEARER回调的相关说明。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 04:00:58