如何使用同一Kafka生产者向关联不同Schema Registry的多Topic生产消息
问题原因说明
首先你当前的配置方式不符合功能设计逻辑:
Kafka生产者的schema.registry.url参数支持填写多个逗号分隔的地址,仅用于同一个Schema Registry集群的高可用容灾,客户端会按顺序遍历地址列表,选择第一个可用的节点处理所有schema校验请求,只有当前节点不可用时才会降级到下一个地址,完全不支持按Topic匹配不同的Schema Registry集群。
你遇到的所有消息都走url2校验的问题,大概率是url1本身不可用,客户端自动降级到了url2节点,才会导致topic1的schema校验失败。
可行解决方案
方案1:创建多个独立生产者实例(推荐)
这是实现最简单、稳定性最高的方案,不同Schema Registry的配置完全隔离,不会互相影响:
- 按照不同的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); }
- 发送消息时根据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
相关产品推荐
相关产品推荐

