如何用Rust连接带认证的Kafka集群并对接Schema Registry?
解决方案
要在Rust中向带认证的Kafka集群(含Schema Registry)生产消息,需要分两部分配置:Kafka生产者的SASL认证,以及Schema Registry的基本认证。由于rdkafka crate仅负责Kafka核心交互,Schema Registry的操作需要依赖专门的crate,以下是完整示例:
1. 依赖配置
在Cargo.toml中添加所需依赖:
[dependencies] rdkafka = { version = "0.34", features = ["ssl", "sasl"] } confluent-schema-registry = "0.12" serde = { version = "1.0", features = ["derive"] } avro_rs = "0.15" tokio = { version = "1.0", features = ["full"] } # 异步生产者需要,同步场景可移除
2. 完整示例代码
use rdkafka::config::ClientConfig; use rdkafka::producer::{FutureProducer, FutureRecord}; use confluent_schema_registry::{Client, Config, BasicAuth}; use serde::Serialize; use std::time::Duration; // 定义要发送的Avro数据结构 #[derive(Serialize, Debug)] struct MyMessage { id: i32, content: String, } #[tokio::main] async fn main() -> Result<(), Box<dyn std::error::Error>> { // 配置Kafka生产者(SASL PLAIN认证) let producer: FutureProducer = ClientConfig::new() .set("bootstrap.servers", "someserver:8080") .set("sasl.mechanism", "PLAIN") .set("security.protocol", "SASL_PLAINTEXT") // 若集群用SSL则改为SASL_SSL .set("sasl.username", "12kjdkjansd") .set("sasl.password", "asd21jaksdqjwdqkwd") .set("acks", "all") .create()?; // 配置Schema Registry客户端(基本认证) let sr_config = Config { url: "https://myrandom.confluent.cloud".to_string(), auth: Some(BasicAuth { username: "asdkjwdnewkeyasdqwd".to_string(), password: "casdkajwcrackthissecretasdkjawd".to_string(), }), ..Default::default() }; let sr_client = Client::new(sr_config)?; // 注册或获取Avro Schema let schema_def = r#" { "type": "record", "name": "MyMessage", "fields": [ {"name": "id", "type": "int"}, {"name": "content", "type": "string"} ] }"#; let schema = sr_client.register_schema("random.topic-value", schema_def, "AVRO").await?; // 序列化消息为Avro字节 let message = MyMessage { id: 1, content: "Hello from Rust!".to_string() }; let avro_bytes = avro_rs::from_value(&schema.schema, &avro_rs::Value::Record( message.into() ))?; // 发送消息到Kafka let record = FutureRecord::to("random.topic") .payload(&avro_bytes) .key("sample-key"); let delivery_result = producer.send(record, Duration::from_secs(10)).await?; println!("消息投递结果: {:?}", delivery_result); Ok(()) }
关键说明
- Kafka认证:通过rdkafka的
sasl.mechanism、sasl.username、sasl.password配置SASL PLAIN认证,security.protocol需匹配集群的安全协议(SASL_PLAINTEXT或SASL_SSL),若用SSL需额外配置ssl.ca.location指定CA证书路径。 - Schema Registry交互:
confluent-schema-registrycrate提供了完整的Schema注册、查询能力,支持基本认证,直接用你提供的SchemaRegistryKey和Secret作为用户名/密码即可。 - 消息序列化:若使用Avro格式,需用
avro_rs或serde_avro将数据结构序列化为符合Schema的字节;若仅发送普通文本/二进制数据,可跳过Schema Registry相关步骤,直接用rdkafka发送原始字节。 - 同步生产者:若不需要异步,可替换
FutureProducer为BaseProducer,移除tokio依赖,调用send方法即可。
内容的提问来源于stack exchange,提问作者William Pham
相关产品推荐
相关产品推荐

