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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 19:54:53