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

