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

如何使用同一Kafka生产者向关联不同Schema Registry的多Topic生产消息

问题原因说明

首先你当前的配置方式不符合功能设计逻辑:
Kafka生产者的schema.registry.url参数支持填写多个逗号分隔的地址,仅用于同一个Schema Registry集群的高可用容灾,客户端会按顺序遍历地址列表,选择第一个可用的节点处理所有schema校验请求,只有当前节点不可用时才会降级到下一个地址,完全不支持按Topic匹配不同的Schema Registry集群。
你遇到的所有消息都走url2校验的问题,大概率是url1本身不可用,客户端自动降级到了url2节点,才会导致topic1的schema校验失败。

可行解决方案

方案1:创建多个独立生产者实例(推荐)

这是实现最简单、稳定性最高的方案,不同Schema Registry的配置完全隔离,不会互相影响:

  1. 按照不同的Schema Registry地址分别初始化生产者Bean
// 对应topic1的生产者,绑定url1
@Bean("producerTopic1")
public Producer producerTopic1() {
    Properties config = sdpProperties();
    config.setProperty("schema.registry.url", "url1");
    config.setProperty("client.id", "producer-topic1");
    // 其余公共配置保持不变
    return new Producer(config);
}

// 对应topic2的生产者,绑定url2
@Bean("producerTopic2")
public Producer producerTopic2() {
    Properties config = sdpProperties();
    config.setProperty("schema.registry.url", "url2");
    config.setProperty("client.id", "producer-topic2");
    // 其余公共配置保持不变
    return new Producer(config);
}
  1. 发送消息时根据Topic选择对应生产者即可
// 注入对应生产者
@Resource(name = "producerTopic1")
private Producer producerTopic1;
@Resource(name = "producerTopic2")
private Producer producerTopic2;

// 发送消息
producerTopic1.send(topic1, genericRecord, datasetId1);
producerTopic2.send(topic2, genericRecord, datasetId2);

方案2:自定义序列化器动态匹配Schema Registry(不推荐)

如果业务场景强制要求只能使用同一个生产者实例,可以重写序列化器的逻辑:在序列化阶段获取当前消息的Topic名称,动态匹配对应的Schema Registry地址,初始化对应客户端完成schema校验和序列化。
该方案需要自行实现多Schema Registry客户端的缓存、Schema缓存、并发安全控制,开发和维护成本较高,没有特殊需求不建议使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 04:06:03