使用MassTransit Kafka与Outbox模式时遇schema.registry.url配置异常
我基于MassTransit搭建了Kafka与Outbox模式,配置代码如下:
var confluentConfiguration = builder.Configuration.GetSection("Kafka").ToConfluentConfiguration(); builder.Services.AddSingleton<ISchemaRegistryClient>(new CachedSchemaRegistryClient(confluentConfiguration)); builder.Services.AddMassTransit(busConfigurator => { busConfigurator.AddEntityFrameworkOutbox<DatabaseContext>(outboxConfigurator => { outboxConfigurator.DuplicateDetectionWindow = outboxOptions.DuplicateDetectionWindow; outboxConfigurator.IsolationLevel = outboxOptions.IsolationLevel; outboxConfigurator.QueryDelay = outboxOptions.QueryDelay; outboxConfigurator.QueryMessageLimit = outboxOptions.QueryMessageLimit; outboxConfigurator.QueryTimeout = outboxOptions.QueryTimeout; outboxConfigurator.UsePostgres(); outboxConfigurator.UseBusOutbox(); }); busConfigurator.UsingInMemory(); // mass transit needs a bus, so we use the in-memory bus busConfigurator.AddRider(rider => { rider.AddProducer<string, UserCreated>("user-created", (context, config) => { config.SetValueSerializer(new ProtobufSerializer<UserCreated>(context.GetRequiredService<ISchemaRegistryClient>())); }); rider.UsingKafka(new ProducerConfig(confluentConfiguration), (context, config) => { config.UseSendFilter(typeof(CustomCloudEventsHeaderFilter<>), context); }); }); });
在Minimal API方法中注入ITopicProducer<string, UserCreated>并调用await userCreatedProducer.Produce(...)时,触发如下异常:
System.InvalidOperationException: No such configuration property: "schema.registry.url" at Confluent.Kafka.Impl.SafeConfigHandle.Set(String name, String value) at Confluent.Kafka.Producer`2.<>c__DisplayClass56_0.<.ctor>b__5(KeyValuePair`2 kvp) at System.Collections.Generic.List`1.ForEach(Action`1 action) at Confluent.Kafka.Producer`2..ctor(ProducerBuilder`2 builder) at Confluent.Kafka.ProducerBuilder`2.Build() at MassTransit.KafkaIntegration.KafkaProducerContext..ctor(ProducerBuilder`2 producerBuilder, IHostConfiguration hostConfiguration, CancellationToken cancellationToken) in /_/src/Transports/MassTransit.KafkaIntegration/KafkaIntegration/KafkaProducerContext.cs:line 24 at MassTransit.KafkaIntegration.ProducerContextFactory.<>c__DisplayClass7_0.<CreateProducer>g__Create|0(ClientContext clientContext, CancellationToken createCancellationToken) in /_/src/Transports/MassTransit.KafkaIntegration/KafkaIntegration/ProducerContextFactory.cs:line 55
已确认schema.registry.url已在confluentConfiguration中配置(若未配置,CachedSchemaRegistryClient构造函数会直接抛出异常)。请问还遗漏了哪些配置?是否需要调用config.Host(...)?
补充说明:appsettings.json配置如下:
"Kafka": { "bootstrap.servers": "localhost:9092", "security.protocol": "SaslPlaintext", "sasl.mechanism": "PLAIN", "sasl.username": "***", "sasl.password": "***", "schema.registry.url": "http://localhost:8081" }
该配置通过ToConfluentConfiguration扩展方法直接读取。
问题根源
你把包含Schema Registry配置的完整Kafka配置直接传给了ProducerConfig,但ProducerConfig仅接受Kafka生产者自身的配置项,schema.registry.url属于Schema Registry客户端的配置,并非生产者配置的一部分,因此Confluent Kafka客户端会抛出配置项不存在的异常。
具体修改方案
方案1:拆分并过滤配置项
从完整配置中过滤掉Schema Registry相关项,再创建ProducerConfig:
var confluentConfiguration = builder.Configuration.GetSection("Kafka").ToConfluentConfiguration(); builder.Services.AddSingleton<ISchemaRegistryClient>(new CachedSchemaRegistryClient(confluentConfiguration)); // 过滤掉Schema Registry相关配置,只保留生产者需要的配置 var producerConfigDict = confluentConfiguration .Where(kv => !kv.Key.StartsWith("schema.registry.")) .ToDictionary(kv => kv.Key, kv => kv.Value); var producerConfig = new ProducerConfig(producerConfigDict); builder.Services.AddMassTransit(busConfigurator => { // ... 保留Outbox相关配置 ... busConfigurator.AddRider(rider => { rider.AddProducer<string, UserCreated>("user-created", (context, config) => { config.SetValueSerializer(new ProtobufSerializer<UserCreated>(context.GetRequiredService<ISchemaRegistryClient>())); }); rider.UsingKafka(producerConfig, (context, config) => { config.UseSendFilter(typeof(CustomCloudEventsHeaderFilter<>), context); }); }); });
方案2:使用MassTransit的config.Host()简化配置
MassTransit的config.Host()方法会自动处理配置项过滤,只加载生产者相关的配置,无需手动拆分:
var confluentConfiguration = builder.Configuration.GetSection("Kafka").ToConfluentConfiguration(); builder.Services.AddSingleton<ISchemaRegistryClient>(new CachedSchemaRegistryClient(confluentConfiguration)); builder.Services.AddMassTransit(busConfigurator => { // ... 保留Outbox相关配置 ... busConfigurator.AddRider(rider => { rider.AddProducer<string, UserCreated>("user-created", (context, config) => { config.SetValueSerializer(new ProtobufSerializer<UserCreated>(context.GetRequiredService<ISchemaRegistryClient>())); }); rider.UsingKafka((context, config) => { // 直接传入配置节,MassTransit会自动过滤非生产者配置 config.Host(builder.Configuration.GetSection("Kafka")); config.UseSendFilter(typeof(CustomCloudEventsHeaderFilter<>), context); }); }); });
方案3:使用配置绑定自动过滤
利用.NET的配置绑定功能,ProducerConfig绑定配置时会自动忽略不匹配的属性:
var producerConfig = new ProducerConfig(); builder.Configuration.GetSection("Kafka").Bind(producerConfig); var confluentConfiguration = builder.Configuration.GetSection("Kafka").ToConfluentConfiguration(); builder.Services.AddSingleton<ISchemaRegistryClient>(new CachedSchemaRegistryClient(confluentConfiguration)); // 后续MassTransit配置中使用上述producerConfig即可
关于config.Host()的说明
是的,调用config.Host()是更推荐的方式,它不仅能自动处理配置项过滤,还能简化配置代码,避免手动处理配置字典的麻烦。
内容的提问来源于stack exchange,提问作者esskar

