Kafka Streams序列化Avro消息时遇Connection reset错误求助
问题:Kafka Streams Avro序列化失败(连接重置)
我之前用Kafka Streams实现了字符串在主题间的流转,现在想改成Avro格式序列化,但一直报错。附上代码和错误信息求助。
我的代码
public static void main(String[] args) throws IOException { String url = "http://url:9092"; Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "TestAvro222"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, url); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); //props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, GenericAvroSerde.class); props.put("schema.registry.url", url); final Map<String, String> serdeConfig = Collections.singletonMap("schema.registry.url", url); final Serde<GenericRecord> valueGenericAvroSerde = new GenericAvroSerde(); valueGenericAvroSerde.configure(serdeConfig, false); StreamsBuilder builder = new StreamsBuilder(); KStream<String, String> textLines = builder.stream ("topic1", Consumed.with(stringSerde, stringSerde)); File file = new File("urltoschema"); Schema schema = new Schema.Parser().parse(file); KTable<String, GenericRecord> textLines_table = textLines.mapValues(value -> avroMaker(schema, value)).toTable(); //KTable textLines_table = textLines.toTable(); textLines_table.toStream().to("TEST-AVRO22", Produced.with(Serdes.String(), valueGenericAvroSerde)); KafkaStreams streams = new KafkaStreams(builder.build(), props); streams.start(); }
报错信息
org.apache.kafka.streams.errors.StreamsException: Error encountered sending record to topic TEST-AVRO22 for task 0_0 due to: org.apache.kafka.common.errors.SerializationException: Error serializing Avro message at org.apache.kafka.streams.processor.internals.RecordCollectorImpl.send(RecordCollectorImpl.java:177) at org.apache.kafka.streams.processor.internals.RecordCollectorImpl.send(RecordCollectorImpl.java:139) at org.apache.kafka.streams.processor.internals.SinkNode.process(SinkNode.java:85) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) at org.apache.kafka.streams.kstream.internals.KStreamMapValues$KStreamMapProcessor.process(KStreamMapValues.java:42) at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) at org.apache.kafka.streams.kstream.internals.KTableSource$KTableSourceProcessor.process(KTableSource.java:152) at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) at org.apache.kafka.streams.kstream.internals.KStreamMapValues$KStreamMapProcessor.process(KStreamMapValues.java:42) at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:146) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forwardInternal(ProcessorContextImpl.java:253) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:232) at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:191) at org.apache.kafka.streams.processor.internals.SourceNode.process(SourceNode.java:84) at org.apache.kafka.streams.processor.internals.StreamTask.lambda$process$1(StreamTask.java:731) at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:809) at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:731) at org.apache.kafka.streams.processor.internals.TaskManager.process(TaskManager.java:1296) at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:784) at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:604) at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:576) Caused by: org.apache.kafka.common.errors.SerializationException: Error serializing Avro message at io.confluent.kafka.serializers.AbstractKafkaAvroSerializer.serializeImpl(AbstractKafkaAvroSerializer.java:101) at io.confluent.kafka.serializers.KafkaAvroSerializer.serialize(KafkaAvroSerializer.java:53) at io.confluent.kafka.streams.serdes.avro.GenericAvroSerializer.serialize(GenericAvroSerializer.java:63) at io.confluent.kafka.streams.serdes.avro.GenericAvroSerializer.serialize(GenericAvroSerializer.java:39) at org.apache.kafka.common.serialization.Serializer.serialize(Serializer.java:62) at org.apache.kafka.streams.processor.internals.RecordCollectorImpl.send(RecordCollectorImpl.java:157) ... 28 common frames omitted Caused by: java.net.SocketException: Connection reset at java.base/java.net.SocketInputStream.read(SocketInputStream.java:186) at java.base/java.net.SocketInputStream.read(SocketInputStream.java:140) at java.base/java.io.BufferedInputStream.fill(BufferedInputStream.java:252) at java.base/java.io.BufferedInputStream.read1(BufferedInputStream.java:292) at java.base/java.io.BufferedInputStream.read(BufferedInputStream.java:351) at java.base/sun.net.www.http.HttpClient.parseHTTPHeader(HttpClient.java:787) at java.base/sun.net.www.http.HttpClient.parseHTTP(HttpClient.java:722) at java.base/sun.net.www.http.HttpClient.parseHTTPHeader(HttpClient.java:896) at java.base/sun.net.www.http.HttpClient.parseHTTP(HttpClient.java:722) at java.base/sun.net.www.protocol.http.HttpURLConnection.getInputStream0(HttpURLConnection.java:1615) at java.base/sun.net.www.protocol.http.HttpURLConnection.getInputStream(HttpURLConnection.java:1520) at java.base/java.net.HttpURLConnection.getResponseCode(HttpURLConnection.java:527) at io.confluent.kafka.schemaregistry.client.rest.RestService.sendHttpRequest(RestService.java:294) at io.confluent.kafka.schemaregistry.client.rest.RestService.httpRequest(RestService.java:384) at io.confluent.kafka.schemaregistry.client.rest.RestService.registerSchema(RestService.java:561) at io.confluent.kafka.schemaregistry.client.rest.RestService.registerSchema(RestService.java:549) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.registerAndGetId(CachedSchemaRegistryClient.java:290) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.register(CachedSchemaRegistryClient.java:397) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.register(CachedSchemaRegistryClient.java:376) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.register(CachedSchemaRegistryClient.java:364) at io.confluent.kafka.schemaregistry.client.SchemaRegistryClient.register(SchemaRegistryClient.java:51) at io.confluent.kafka.serializers.AbstractKafkaAvroSerializer.serializeImpl(AbstractKafkaAvroSerializer.java:70) ... 33 common frames omitted
解决方案
1. 修正Schema Registry地址
你把Kafka Broker的地址(http://url:9092)错误配置给了Schema Registry。Kafka Broker使用TCP协议,默认端口9092;而Schema Registry是HTTP服务,默认端口8081。你需要将schema.registry.url改为正确的Schema Registry地址,比如:
String schemaRegistryUrl = "http://your-schema-registry-host:8081"; props.put("schema.registry.url", schemaRegistryUrl); // 同时更新Serde配置里的地址 final Map<String, String> serdeConfig = Collections.singletonMap("schema.registry.url", schemaRegistryUrl);
2. 解决stringSerde未初始化问题
代码中Consumed.with(stringSerde, stringSerde)里的stringSerde没有定义,需要提前初始化:
Serde<String> stringSerde = Serdes.String();
3. 验证Schema Registry可用性
确认Schema Registry服务正在运行,并且你的应用服务器能访问到该地址(可以用curl命令测试:curl http://your-schema-registry-host:8081/subjects)。
4. 检查avroMaker方法
确保avroMaker方法生成的GenericRecord完全符合你加载的Schema定义,字段名、类型都要匹配,否则即使连接Schema Registry成功,也会出现序列化错误。
内容的提问来源于stack exchange,提问作者yaboi991
相关产品推荐
相关产品推荐

