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

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>

问题根源分析

  1. 聚合类型不匹配:聚合初始值用了AccountHolder类型,但聚合器返回ProcessedStream类型,Kafka Streams无法强制转换,直接触发类转换异常。
  2. Serde实例错误:processedStreamSerde错误复用了AccountHolder的Serde实例,未使用对应ProcessedStream的Serde。
  3. 状态存储类型不匹配:Materialized指定的value Serde是AccountHolder的,但实际聚合状态是ProcessedStream,导致序列化失败。
  4. 计数逻辑错误:每次聚合都重置count为1再加1,永远得到2,未正确累计交易次数。
  5. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 14:04:58