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

Kafka:为何groupBy与reduce操作需配置默认Key和Value Serde?

Kafka Streams + Quarkus:groupBy/reduce阶段缺失Key/Value Serde配置报错

我用Quarkus开发流式应用,处理流程如下:

  • 使用flatMap修改Key并将单条消息拆分为多条;
  • 利用步骤1的Key与KTable执行join操作;
  • 通过transform实现有状态操作;
  • 再次使用flatMap将Key改回步骤1之前的原始值;
  • 以步骤4的Key执行groupBy操作;
  • 通过reduce将记录合并为包含JSON数组的单条消息。

实际效果是:把Key为id1的入站消息拆分为Key为k1、k2等的多条消息,经join和transform增强后,将Key改回id1,最终合并为单条Key为id1的消息。

但执行第5、6步时,一直收到需要配置默认Key和Value Serde的错误,跳过这两步应用就能正常运行。

异常信息

2022-10-17 16:42:34,884 ERROR [org.apa.kaf.str.KafkaStreams] (app-alerts-6a7c4df8-7813-4d5d-9a86-d6f3db7c8ef0-StreamThread-1) stream-client [app-alerts-6a7c4df8-7813-4d5d-9a86-d6f3db7c8ef0] Encountered the following exception during processing and the registered exception handler opted to SHUTDOWN_CLIENT. The streams client is going to shut down now. : org.apache.kafka.streams.errors.StreamsException: org.apache.kafka.common.config.ConfigException: Please specify a key serde or set one through StreamsConfig#DEFAULT_KEY_SERDE_CLASS_CONFIG
        at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:627)
        at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:551)
Caused by: org.apache.kafka.common.config.ConfigException: Please specify a key serde or set one through StreamsConfig#DEFAULT_KEY_SERDE_CLASS_CONFIG
        at org.apache.kafka.streams.StreamsConfig.defaultKeySerde(StreamsConfig.java:1587)
        at org.apache.kafka.streams.processor.internals.AbstractProcessorContext.keySerde(AbstractProcessorContext.java:90)
        at org.apache.kafka.streams.processor.internals.SerdeGetter.keySerde(SerdeGetter.java:47)
        at org.apache.kafka.streams.kstream.internals.WrappingNullableUtils.prepareSerde(WrappingNullableUtils.java:63)
        at org.apache.kafka.streams.kstream.internals.WrappingNullableUtils.prepareKeySerde(WrappingNullableUtils.java:90)
        at org.apache.kafka.streams.state.internals.MeteredKeyValueStore.initStoreSerde(MeteredKeyValueStore.java:195)
        at org.apache.kafka.streams.state.internals.MeteredKeyValueStore.init(MeteredKeyValueStore.java:144)
        at org.apache.kafka.streams.processor.internals.ProcessorStateManager.registerStateStores(ProcessorStateManager.java:212)
        at org.apache.kafka.streams.processor.internals.StateManagerUtil.registerStateStores(StateManagerUtil.java:97)
        at org.apache.kafka.streams.processor.internals.StreamTask.initializeIfNeeded(StreamTask.java:231)
        at org.apache.kafka.streams.processor.internals.TaskManager.tryToCompleteRestoration(TaskManager.java:454)
        at org.apache.kafka.streams.processor.internals.StreamThread.initializeAndRestorePhase(StreamThread.java:865)
        at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:747)
        at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:589)
        ... 1 more

当前StreamsConfig配置

acceptable.recovery.lag = 10000
        application.id = machine-alerts
        application.server = 
        bootstrap.servers = [kafka:9092]
        buffered.records.per.partition = 1000
        built.in.metrics.version = latest
        cache.max.bytes.buffering = 10240
        client.id = 
        commit.interval.ms = 1000
        connections.max.idle.ms = 540000
        default.deserialization.exception.handler = class org.apache.kafka.streams.errors.LogAndFailExceptionHandler
        default.dsl.store = rocksDB
        default.key.serde = null
        default.list.key.serde.inner = null
        default.list.key.serde.type = null
        default.list.value.serde.inner = null
        default.list.value.serde.type = null
        default.production.exception.handler = class org.apache.kafka.streams.errors.DefaultProductionExceptionHandler
        default.timestamp.extractor = class org.apache.kafka.streams.processor.FailOnInvalidTimestamp
        default.value.serde = null
        max.task.idle.ms = 0
        max.warmup.replicas = 2
        metadata.max.age.ms = 500
        metric.reporters = []
        metrics.num.samples = 2
        metrics.recording.level = DEBUG
        metrics.sample.window.ms = 30000
        num.standby.replicas = 0
        num.stream.threads = 1
        poll.ms = 100
        probing.rebalance.interval.ms = 600000
        processing.guarantee = at_least_once
        rack.aware.assignment.tags = []
        receive.buffer.bytes = 32768
        reconnect.backoff.max.ms = 1000
        reconnect.backoff.ms = 50
        repartition.purge.interval.ms = 30000
        replication.factor = -1
        request.timeout.ms = 40000
        retries = 0
        retry.backoff.ms = 100
        rocksdb.config.setter = null
        security.protocol = PLAINTEXT
        send.buffer.bytes = 131072
        state.cleanup.delay.ms = 600000
        state.dir = /tmp/kafka-streams
        task.timeout.ms = 300000
        topology.optimization = none
        upgrade.from = null
        window.size.ms = null
        windowed.inner.class.serde = null
        windowstore.changelog.additional.retention.ms = 86400000

解决方案

从配置能看到default.key.serde和default.value.serde都是null,这就是问题根源。Kafka Streams在执行groupBy和reduce这类需要状态存储的操作时,必须明确知道Key和Value的序列化/反序列化方式,要么全局配置默认Serde,要么在操作时指定。

方法1:全局配置默认Serde

在Quarkus的配置文件(比如application.properties)中添加:

# 假设你的Key是字符串类型,Value是JSON格式
kafka-streams.default.key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
kafka-streams.default.value.serde=org.apache.kafka.common.serialization.Serdes$StringSerde
# 如果Value是自定义对象,用Quarkus自带的JSON Serde:
# kafka-streams.default.value.serde=io.quarkus.kafka.streams.runtime.serialization.JsonSerde

方法2:在groupBy/reduce时指定Serde

如果不想全局配置,可以在调用groupBy和reduce时显式指定Serde:

// 假设Key是String,Value是你的自定义类型或String
stream.groupBy(
    (key, value) -> key, // 你的Key提取逻辑
    Grouped.with(Serdes.String(), Serdes.String()) // 指定Key和Value的Serde
)
.reduce(
    (value1, value2) -> mergeValuesIntoJsonArray(value1, value2), // 你的合并逻辑
    Materialized.with(Serdes.String(), Serdes.String()) // 状态存储的Serde
);

为什么前几步没问题?

前几步的flatMap、join、transform操作,可能因为你在数据源(比如KStream)创建时已经指定了Serde,或者Quarkus自动推断了Serde,但groupBy会触发状态存储的创建,此时必须有明确的Serde配置,否则就会抛出这个错误。


内容的提问来源于stack exchange,提问作者Chain Head

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 07:21:36