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

如何用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-registry crate提供了完整的Schema注册、查询能力,支持基本认证,直接用你提供的SchemaRegistryKey和Secret作为用户名/密码即可。
  • 消息序列化:若使用Avro格式,需用avro_rs或serde_avro将数据结构序列化为符合Schema的字节;若仅发送普通文本/二进制数据,可跳过Schema Registry相关步骤,直接用rdkafka发送原始字节。
  • 同步生产者:若不需要异步,可替换FutureProducer为BaseProducer,移除tokio依赖,调用send方法即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 04:23:16