Kafka Streams聚合Protobuf消息时类转换异常求助
Kafka Streams Protobuf聚合类转换异常问题解决
我有一条基于AccountHolder.proto(包含name、amount、time字段)的Protobuf消息流,需要对数据进行聚合操作后,将结果以ProcessedStream.proto(包含count、balance、time字段)格式的Protobuf消息输出。但在StreamOperation.java的aggregate()方法处抛出类转换异常,相关代码如下,希望能实现按用户聚合总余额并映射为ProcessedStream格式。
原问题代码
StreamOperation.java
KafkaProtobufSerde<AccountHolderOuterClass.AccountHolder> accountHolderSerde = serde.accHolderSerde(props); KafkaProtobufSerde<AccountHolderOuterClass.AccountHolder> processedStreamSerde = serde.accHolderSerde(props); StreamsBuilder builder = new StreamsBuilder(); KStream<String, AccountHolderOuterClass.AccountHolder> stream = builder.stream("streams-test-app-input", Consumed.with(Serdes.String(),accountHolderSerde)); AccountHolderOuterClass.AccountHolder initialValues = AccountHolderOuterClass.AccountHolder.newBuilder() .setAmount(0).setName(null).setTime(Instant.ofEpochMilli(0L).toString()) .build(); KTable<String, ProcessedStreamOuterClass.ProcessedStream> bankbalance = stream.groupByKey(Grouped.with(Serdes.String(),accountHolderSerde)) .aggregate(() -> initialValues,(key, transaction, balance) -> balanceCalc(transaction, balance), Named.as("balance-agg"),Materialized.with(Serdes.String(),accountHolderSerde)); bankbalance.toStream().to("bank-trans-output",Produced.with(Serdes.String(),processedStreamSerde)); KafkaStreams streams = new KafkaStreams(builder.build(),props); streams.start(); } private static ProcessedStreamOuterClass.ProcessedStream balanceCalc(AccountHolderOuterClass.AccountHolder transaction, AccountHolderOuterClass.AccountHolder balance){ Long bal_time = Instant.parse(balance.getTime()).toEpochMilli(); Long trans_time = Instant.parse(transaction.getTime()).toEpochMilli(); Instant max = Instant.ofEpochMilli(Math.max(bal_time,trans_time)); int count = 1; count++; ProcessedStreamOuterClass.ProcessedStream totalBalance = ProcessedStreamOuterClass.ProcessedStream.newBuilder() .setBalance(balance.getAmount()+transaction.getAmount()) .setCount(Integer.toString(count)) .setTime(max.toString()) .build(); return totalBalance; }
AccountHolder.proto
syntax = "proto3"; package holder; option java_package = "dk.danske.model.accholder"; message AccountHolder{ string name=1; int32 amount=2; string time=3; }
ProtobufSerde.java
public KafkaProtobufSerde<AccountHolderOuterClass.AccountHolder> accHolderSerde(Properties envProps) { final KafkaProtobufSerde<AccountHolderOuterClass.AccountHolder> protobufSerde = new KafkaProtobufSerde<>(); Map<String, String> serdeConfig = new HashMap<>(); serdeConfig.put(SCHEMA_REGISTRY_URL_CONFIG, envProps.getProperty("schema.registry.url")); protobufSerde.configure(serdeConfig, false); return protobufSerde; } public KafkaProtobufSerde<ProcessedStreamOuterClass.ProcessedStream> processedStreamSerde(Properties envProps) { final KafkaProtobufSerde<ProcessedStreamOuterClass.ProcessedStream> protobufSerde = new KafkaProtobufSerde<>(); Map<String, String> serdeConfig = new HashMap<>(); serdeConfig.put(SCHEMA_REGISTRY_URL_CONFIG, envProps.getProperty("schema.registry.url")); protobufSerde.configure(serdeConfig, false); return protobufSerde; }
KafkaProducer.java
Properties props = new Properties(); props.put("bootstrap.servers", "http://127.0.0.1:9200"); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "io.confluent.kafka.serializers.protobuf.KafkaProtobufSerializer"); props.put("schema.registry.url", "http://localhost:3200"); KafkaProducer<String, AccountHolderOuterClass.AccountHolder> producer = new KafkaProducer<>(props); List<String> nameList = new ArrayList<>(); nameList.add("Johny"); nameList.add("Thomas"); nameList.add("Roger"); nameList.add("Micheal"); nameList.add("William"); nameList.add("Jonas"); nameList.add("Mike"); nameList.add("Kim"); for( String name:nameList){ producer.send(createTransaction(name)); Thread.sleep(100); log.info("proto message sent to kafka producer",createTransaction(name)); } producer.flush(); producer.close(); } private static ProducerRecord<String, AccountHolderOuterClass.AccountHolder> createTransaction(String name) { AccountHolderOuterClass.AccountHolder accountHolder = AccountHolderOuterClass.AccountHolder.newBuilder() .setAmount(ThreadLocalRandom.current().nextInt(0,9000)) .setName(name) .setTime(Instant.now().toString()) .build(); return new ProducerRecord<>(TOPIC,name,accountHolder); }
pom.xml
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-streams</artifactId> <version>3.5.1</version> </dependency> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-streams-protobuf-serde</artifactId> <version>7.5.1</version> </dependency> <dependency> <groupId>org.slf4j</groupId> <artifactId>slf4j-api</artifactId> <version>2.0.9</version> </dependency> <dependency> <groupId>org.slf4j</groupId> <artifactId>slf4j-log4j12</artifactId> <version>2.0.9</version> </dependency> <dependency> <groupId>com.google.protobuf</groupId> <artifactId>protobuf-java</artifactId> <version>3.24.0</version> </dependency> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-protobuf-serializer</artifactId> <version>5.5.1</version> </dependency> </dependencies> <build> <plugins> <plugin> <groupId>org.apache.maven.plugins</groupId> <artifactId>maven-compiler-plugin</artifactId> <version>3.8.1</version> <configuration> <source>17</source> <target>17</target> </configuration> </plugin> <plugin> <groupId>com.github.os72</groupId> <artifactId>protoc-jar-maven-plugin</artifactId> <version>3.1.0.1</version> <executions> <execution> <phase>generate-sources</phase> <goals> <goal>run</goal> </goals> <configuration> <protocVersion>3.1.0</protocVersion> <inputDirectories> <include>src/main/resources</include> </inputDirectories> </configuration> </execution> </executions> </plugin> </plugins> </build>
问题根源分析
- 聚合类型不匹配:聚合初始值用了
AccountHolder类型,但聚合器返回ProcessedStream类型,Kafka Streams无法强制转换,直接触发类转换异常。 - Serde实例错误:
processedStreamSerde错误复用了AccountHolder的Serde实例,未使用对应ProcessedStream的Serde。 - 状态存储类型不匹配:
Materialized指定的value Serde是AccountHolder的,但实际聚合状态是ProcessedStream,导致序列化失败。 - 计数逻辑错误:每次聚合都重置count为1再加1,永远得到2,未正确累计交易次数。
- Protobuf字段不合理:
ProcessedStream的count字段定义为string,计数场景用int32类型更合适。
修正后的解决方案
1. 补充ProcessedStream.proto定义
syntax = "proto3"; package processed; option java_package = "dk.danske.model.processed"; message ProcessedStream { int32 count = 1; // 改为int32类型适配计数场景 int32 balance = 2; string time = 3; }
2. 修正StreamOperation.java
// 获取正确的Serde实例 KafkaProtobufSerde<AccountHolderOuterClass.AccountHolder> accountHolderSerde = serde.accHolderSerde(props); KafkaProtobufSerde<ProcessedStreamOuterClass.ProcessedStream> processedStreamSerde = serde.processedStreamSerde(props); StreamsBuilder builder = new StreamsBuilder(); KStream<String, AccountHolderOuterClass.AccountHolder> stream = builder.stream( "streams-test-app-input", Consumed.with(Serdes.String(), accountHolderSerde) ); // 初始值改为ProcessedStream类型,匹配聚合结果类型 ProcessedStreamOuterClass.ProcessedStream initialValues = ProcessedStreamOuterClass.ProcessedStream.newBuilder() .setBalance(0) .setCount(0) .setTime(Instant.ofEpochMilli(0L).toString()) .build(); // 修正聚合操作的类型匹配 KTable<String, ProcessedStreamOuterClass.ProcessedStream> bankbalance = stream.groupByKey(Grouped.with(Serdes.String(), accountHolderSerde)) .aggregate( () -> initialValues, (key, transaction, currentState) -> balanceCalc(transaction, currentState), Named.as("balance-agg"), // 状态存储使用ProcessedStream的Serde Materialized.with(Serdes.String(), processedStreamSerde) ); // 输出时使用正确的Serde bankbalance.toStream().to( "bank-trans-output", Produced.with(Serdes.String(), processedStreamSerde) ); KafkaStreams streams = new KafkaStreams(builder.build(), props); streams.start(); } // 修正参数和返回类型,正确累计交易次数和余额 private static ProcessedStreamOuterClass.ProcessedStream balanceCalc( AccountHolderOuterClass.AccountHolder transaction, ProcessedStreamOuterClass.ProcessedStream currentState ) { Long stateTime = Instant.parse(currentState.getTime()).toEpochMilli(); Long transTime = Instant.parse(transaction.getTime()).toEpochMilli(); Instant maxTime = Instant.ofEpochMilli(Math.max(stateTime, transTime)); // 累计交易次数和余额 int newCount = currentState.getCount() + 1; int newBalance = currentState.getBalance() + transaction.getAmount(); return ProcessedStreamOuterClass.ProcessedStream.newBuilder() .setBalance(newBalance) .setCount(newCount) .setTime(maxTime.toString()) .build(); }
3. 其他关键修正
- 生产者地址修正:KafkaProducer.java中
bootstrap.servers改为Kafka默认端口(通常是9092),而非Elasticsearch的9200:props.put("bootstrap.servers", "localhost:9092"); - 依赖版本统一:pom.xml中
kafka-protobuf-serializer版本改为和kafka-streams-protobuf-serde一致的7.5.1,避免版本冲突。 - protoc版本升级:将maven插件的
protocVersion改为3.24.0,和protobuf-java版本匹配。
内容的提问来源于stack exchange,提问作者Jyothi Siva
相关产品推荐
相关产品推荐

