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
相关产品推荐
相关产品推荐

