使用Schema Registry的C# Kafka生产者报错:无效接收大小
问题描述
本地使用C#学习Kafka,此前基于字符串的生产者、消费者与流处理均正常运行,现在尝试通过Schema Registry处理复杂类型。编写控制台程序注册自定义类的Schema时,运行抛出HttpRequestException,Kafka容器日志显示InvalidReceiveException(提示接收大小超过100MB限制),但实际消息体积很小。相关代码与配置如下:
注册Schema的代码
static async Task Main(string[] args) { var schemaRegistryConfig = new SchemaRegistryConfig { Url = "http://localhost:29092" }; var schemaRegistry = new CachedSchemaRegistryClient(schemaRegistryConfig); var schema = @"{ ""type"": ""TotallyCoolCustomClass"", ""properties"": { ""FavouriteQuote"": {""type"": ""string""}, ""FavouriteNumber"": {""type"": ""integer""} } }"; var subject = "SimpleTest"; var schemaId = await schemaRegistry.RegisterSchemaAsync(subject, schema); }
错误信息
System.Net.Http.HttpRequestException: '[http://localhost:29092/] HttpRequestException: An error occurred while sending the request.'
Kafka容器日志
WARN [SocketServer listenerType=ZK_BROKER, nodeId=1] Unexpected error from /172.22.0.1 (channelId=172.22.0.3:29092-172.22.0.1:34802-28); closing connection (org.apache.kafka.common.network.Selector) org.apache.kafka.common.network.InvalidReceiveException: Invalid receive (size = 1347375956 larger than 104857600) at org.apache.kafka.common.network.NetworkReceive.readFrom(NetworkReceive.java:94)
Docker Compose配置
version: '2' services: zookeeper: image: confluentinc/cp-zookeeper:latest environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 ports: - 22181:2181 kafka: image: confluentinc/cp-kafka:latest depends_on: - zookeeper ports: - 29092:29092 environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: PLAINTEXT KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 schema-registry: image: confluentinc/cp-schema-registry:latest depends_on: - kafka ports: - 28081:28081 environment: SCHEMA_REGISTRY_KAFKASTORE_CONNECTION_URL: 'zookeeper:2181' SCHEMA_REGISTRY_HOST_NAME: schema-registry SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: 'PLAINTEXT://kafka:9092' # SCHEMA_REGISTRY_LISTENERS: http://0.0.0.0:28081
生产者代码(参考)
public class KafkaPProducerHostedService : IHostedService { private readonly ILogger<KafkaPProducerHostedService> _logger; private readonly IProducer<string, TotallyCoolCustomClass> _producer; public KafkaPProducerHostedService(ILogger<KafkaPProducerHostedService> logger) { _logger = logger; var config = new ProducerConfig(); config.BootstrapServers = "localhost:29092"; var schemaRegistryConfig = new SchemaRegistryConfig { Url = "http://localhost:28081/" }; //video on serializers _producer = new ProducerBuilder<string, TotallyCoolCustomClass>(config) .SetValueSerializer(new JsonSerializer<TotallyCoolCustomClass>(new CachedSchemaRegistryClient(schemaRegistryConfig))) .Build(); } public async Task StartAsync(CancellationToken cancellationToken) { int i = 0; var rand = new Random(); while (true) { Thread.Sleep(2000); var customClass = new TotallyCoolCustomClass(rand); var key = $"UpdatedKey-{i}"; await _producer.ProduceAsync("SimpleTest", new Message<string, TotallyCoolCustomClass>() { Key = key, Value = customClass }, cancellationToken); Console.WriteLine($"Published: Key: {key} with favourite quote {customClass.FavouriteQuote} and favourite number {customClass.FavouriteNumber}"); i++; } } public Task StopAsync(CancellationToken cancellationToken) { _producer?.Dispose(); return Task.CompletedTask; } }
问题原因与解决方案
核心问题
注册Schema的代码中,Schema Registry的URL配置错误:你将Url设置成了Kafka Broker的端口29092,而不是Schema Registry的端口28081。
当用HTTP请求访问Kafka Broker的端口时,Broker会尝试解析Kafka协议的数据包,但收到的是HTTP请求内容,解析后会得到一个异常大的"消息大小"(HTTP头的二进制内容被当成了Kafka协议的消息长度字段),这就触发了InvalidReceiveException,提示接收大小超过100MB限制。
修复步骤
修改注册Schema代码中的SchemaRegistryConfig,将URL改为Schema Registry的端口:
var schemaRegistryConfig = new SchemaRegistryConfig { Url = "http://localhost:28081" };
重启控制台程序,重新执行注册操作即可正常完成Schema注册。
额外验证
从你提供的生产者代码可以看到,生产者部分已经正确配置了Schema Registry的URL为http://localhost:28081/,说明你清楚两者的端口区别,只是在注册Schema的控制台程序中疏忽了配置。
内容的提问来源于Stack Exchange,提问作者Joshua Mee
相关产品推荐
相关产品推荐

