Apache Kafka Streams 数据类型与序列化(Serdes)完全指南
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
导读
Kafka Streams 作为一个基于 Kafka 的流处理库,其所有数据都以字节流的形式在 Broker 与客户端之间传输,因此每个 Kafka Streams 应用都必须为记录键(Key)和记录值(Value)提供 Serde(Serializer/Deserializer,序列化器/反序列化器)。本文以官方开发者指南 datatypes.md 为骨架,结合 Kafka 仓库源码,系统讲解 Serde 的配置方式、内置 Serde 清单、窗口 Serde(Windowed Serdes)、自定义 Serde 实现以及 Scala DSL 的隐式 Serde 机制。读完后你将能熟练地在自己的 Streams 应用中配置、覆盖、组合与定制 Serde,并理解其底层实现原理。
说明:
docs/documentation/streams/developer-guide/datatypes.md为文档站点的重定向占位页,真实内容位于 docs/streams/developer-guide/datatypes.md,本文以其真实内容为准。
为什么每个 Kafka Streams 应用都必须提供 Serdes
Kafka 本身只关心字节数组:Producer 把任意类型序列化成byte[]写入 Topic,Consumer 再把byte[]反序列化回业务对象。Kafka Streams 运行在这层字节语义之上,因此应用必须在需要物化(materialize)数据时,为记录的键和值指明“如何从对象到字节、从字节回对象”。
从源码结构看,需要 Serde 信息的操作包括:stream()、table()、to()、repartition()、groupByKey()、groupBy()等 DSL 方法。这些方法要么直接消费外部 Topic 的字节数据,要么把内部计算结果写回新的 Topic 或状态存储,任何一处缺少序列化能力都会导致运行时失败。
提供 Serde 有且仅有两种方式(至少使用其中一种):
- 在
java.util.Properties配置实例中设置默认 Serde(通过StreamsConfig的两个配置项); - 在调用相应 API 方法时显式传入 Serde,从而覆盖默认配置。
配置默认 Serdes
在 Kafka Streams 配置中指定的 Serde 将作为整个应用的默认值。由于该配置项的默认值是null,你必须通过配置设置默认 Serde,或按下文介绍的方式显式传参,二者必居其一。
import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.StreamsConfig; Properties settings = new Properties(); // 记录键的默认 Serde(此处为 String 类型的内置 Serde) settings.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); // 记录值的默认 Serde(此处为 Long 类型的内置 Serde) settings.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.Long().getClass().getName());两个关键配置项:
StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG(配置名default.key.serde):键的默认 Serde 类;StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG(配置名default.value.serde):值的默认 Serde 类。
注意这里写入的是类的全限定名(Class.getName()),因为 Streams 会通过反射Utils.newInstance(...)实例化该类,这就要求配置指向的 Serde 类必须拥有无参构造函数。这也是本文后面“自定义 Serde”一节强调“自定义 Serde 必须是无泛型的具体类”的根本原因。
覆盖默认 Serdes(显式传参)
你可以在调用 API 方法时显式传入 Serde 来覆盖默认设置:
import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serdes; final Serde<String> stringSerde = Serdes.String(); final Serde<Long> longSerde = Serdes.Long(); // KStream userCountByRegion 的键是 String(区域),值是 Long(用户数) KStream<String, Long> userCountByRegion = ...; userCountByRegion.to("RegionCountsTopic", Produced.with(stringSerde, longSerde));DSL 中负责承载 Serde 参数的是Consumed、Produced、Grouped、Joined、Materialized等“伴随类”,它们都有静态工厂方法与with(...)组合方式。
如果你只想选择性覆盖——保留部分字段使用默认 Serde,那么对应位置不传即可:
import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serdes; // 键保持默认 String Serde(不指定),只覆盖值的默认 Serde 为 Long final Serde<Long> longSerde = Serdes.Long(); KStream<String, Long> userCountByRegion = ...; userCountByRegion.to("RegionCountsTopic", Produced.valueSerde(Serdes.Long()));反序列化异常处理
如果部分流入的记录损坏或格式不正确,反序列化器会抛出异常。自 1.0.x 起,Kafka 引入了DeserializationExceptionHandler接口,允许你自定义对这类记录的处理策略(例如丢弃坏记录并继续、或终止应用)。自定义实现通过StreamsConfig指定,详细配置见 Configuring a Streams Application 中的deserialization.exception.handler一节(旧的default.deserialization.exception.handler配置名已弃用,具体说明可参见 config-streams.md)。
内置 Serdes
基本类型与常用类型
Kafka 在kafka-clientsMaven 构件中为 Java 基本类型及byte[]等提供了大量内置 Serde 实现:
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-clients</artifactId> <version>4.3.0</version> </dependency>这些实现位于org.apache.kafka.common.serialization包(源码见 clients/src/main/java/org/apache/kafka/common/serialization),可通过Serdes工厂类获取:
| 数据类型 | Serde |
|---|---|
byte[] | Serdes.ByteArray()、Serdes.Bytes()(见下方提示) |
ByteBuffer | Serdes.ByteBuffer() |
Double | Serdes.Double() |
Integer | Serdes.Integer() |
Long | Serdes.Long() |
String | Serdes.String() |
UUID | Serdes.UUID() |
Void | Serdes.Void() |
List | Serdes.ListSerde() |
Boolean | Serdes.Boolean() |
提示:
Bytes是 Javabyte[]的包装类(见 clients/src/main/java/org/apache/kafka/common/utils/Bytes.java),提供了正确的equals与排序语义。相比裸byte[],在应用中使用Bytes更安全、更推荐。
从源码看(Serdes.java),这些工厂方法背后是一一对应的内部静态类,例如LongSerde内部就是new WrapperSerde<>(new LongSerializer(), new LongDeserializer())。WrapperSerde是Serde接口的一个便捷实现,把configure/close/serializer/deserializer四个方法全部委托给内部的 Serializer 与 Deserializer,这也解释了为什么大多数内置 Serde 本身不包含业务逻辑——它们只是“序列化器 + 反序列化器”的组合器。
此外,Serdes.serdeFrom(Class<T> type)可以根据类型自动推断内置 Serde(支持 String、Short、Integer、Long、Float、Double、byte[]、ByteBuffer、Bytes、UUID、Boolean),遇到未知类型会抛出IllegalArgumentException;而Serdes.serdeFrom(Serializer<T>, Deserializer<T>)则可以从独立的序列化器与反序列化器组合出任意 Serde。
ListSerde值得一提:它有两种构造方式——无参Serdes.ListSerde()直接使用内置的ListSerializer/ListDeserializer,以及带参数Serdes.ListSerde(Class<L> listClass, Serde<Inner> innerSerde),后者允许指定 List 的具体实现类(如ArrayList.class)与元素类型的内置 Serde,从而支持泛型列表的序列化。
JSON
Kafka Streams 官方代码示例中包含一个基础的 JSON Serde 实现:
- PageViewTypedDemo.java
如示例所示,可以通过Serdes.serdeFrom(<serializer实例>, <deserializer实例>)组合出 JSON 兼容的序列化器与反序列化器。示例中的JSONSerde<T>类同时实现了Serializer<T>、Deserializer<T>与Serde<T>三个接口,内部基于 Jackson 的ObjectMapper:反序列化时用@JsonTypeInfo注解记录的_t字段区分具体子类型,配合@JsonSubTypes注册所有可能的类型。这种“泛型 JSON Serde + 类型字段”的模式非常适合一类数据对应多种具体结构的场景。
窗口 Serdes(Windowed Serdes)
时间窗口(Tumbling、Hopping、Session Window)是 Kafka Streams 最核心的抽象之一。窗口操作(如windowedBy、groupByKey+ 窗口聚合)产生的键类型是Windowed<K>——它把“原始键 + 窗口起始/结束时间”打包在一起。kafka-streams构件为这类窗口类型提供了专门的 Serde 实现:
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-streams</artifactId> <version>4.3.0</version> </dependency>这些实现位于org.apache.kafka.streams.kstream包(源码见 streams/src/main/java/org/apache/kafka/streams/kstream):
Serdes(组合类):
WindowedSerdes.TimeWindowedSerde<T>WindowedSerdes.SessionWindowedSerde<T>
序列化器:
TimeWindowedSerializer<T>SessionWindowedSerializer<T>
反序列化器:
TimeWindowedDeserializer<T>SessionWindowedDeserializer<T>
从源码看,TimeWindowedSerializer.serialize()最终委托给 WindowKeySchema.toBinary(...) 完成字节编码;TimeWindowedDeserializer.deserialize()则根据是否 changelog Topic 分别调用WindowKeySchema.fromStoreKey(...)或WindowKeySchema.from(...)。因此窗口 Serde 的底层字节布局与状态存储(State Store)的键模式完全一致,这也保证了窗口聚合结果能无缝落盘与恢复。
代码中的用法
// 时间窗口 Serde —— 工厂方法 Serde<Windowed<String>> timeWindowedSerde = WindowedSerdes.timeWindowedSerdeFrom(String.class, 500L); // 时间窗口 Serde —— 构造函数 Serde<Windowed<String>> timeWindowedSerde2 = new WindowedSerdes.TimeWindowedSerde<>(Serdes.String(), 500L); // 会话窗口 Serde —— 工厂方法 Serde<Windowed<String>> sessionWindowedSerde = WindowedSerdes.sessionWindowedSerdeFrom(String.class); // 会话窗口 Serde —— 构造函数 Serde<Windowed<String>> sessionWindowedSerde2 = new WindowedSerdes.SessionWindowedSerde<>(Serdes.String()); // 单独使用序列化器 / 反序列化器 TimeWindowedSerializer<String> serializer = new TimeWindowedSerializer<>(Serdes.String().serializer()); TimeWindowedDeserializer<String> deserializer = new TimeWindowedDeserializer<>(Serdes.String().deserializer(), 500L);窗口大小(如上面示例中的500L,单位毫秒)对时间窗口 Serde 至关重要:反序列化器需要它来计算窗口的结束时间(窗口的字节编码中通常只保存起始时间戳)。会话窗口没有固定大小,因此不需要该参数。
命令行工具中的用法
使用命令行工具(如bin/kafka-console-consumer.sh)消费窗口结果 Topic 时,可以通过--formatter-property传入窗口反序列化器与窗口大小,属性名遵循前缀模式:
# 时间窗口反序列化器配置 --formatter-property print.key=true \ --formatter-property key.deserializer=org.apache.kafka.streams.kstream.TimeWindowedDeserializer \ --formatter-property key.deserializer.windowed.inner.deserializer.class=org.apache.kafka.common.serialization.StringDeserializer \ --formatter-property key.deserializer.window.size.ms=500 # 会话窗口反序列化器配置 --formatter-property print.key=true \ --formatter-property key.deserializer=org.apache.kafka.streams.kstream.SessionWindowedDeserializer \ --formatter-property key.deserializer.windowed.inner.deserializer.class=org.apache.kafka.common.serialization.StringDeserializer这里的key.deserializer指向无参构造的TimeWindowedDeserializer/SessionWindowedDeserializer,windowed.inner.deserializer.class指定窗口内部键的反序列化器,window.size.ms指定窗口大小(毫秒)。会话窗口无需window.size.ms。
已弃用的配置
以下StreamsConfig参数已弃用,应改用向序列化器/反序列化器构造函数直接传参的方式:
StreamsConfig.WINDOWED_INNER_CLASS_SERDE→ 改用TimeWindowedSerializer.WINDOWED_INNER_SERIALIZER_CLASS(windowed.inner.serializer.class)与TimeWindowedDeserializer.WINDOWED_INNER_DESERIALIZER_CLASS(windowed.inner.deserializer.class)StreamsConfig.WINDOW_SIZE_MS_CONFIG→ 改用TimeWindowedDeserializer.WINDOW_SIZE_MS_CONFIG(window.size.ms)
从 TimeWindowedSerializer.java 的configure()实现可以看出兼容逻辑:它优先读取windowed.inner.serializer.class,若为空则回退读取已弃用的StreamsConfig.WINDOWED_INNER_CLASS_SERDE并打印弃用告警;同时校验“构造函数传入的 inner 与配置指定的 inner”必须一致,否则抛出IllegalArgumentException。TimeWindowedDeserializer的窗口大小解析也有类似的约束:构造函数与window.size.ms配置只能二选一,二者都设置或都不设置都会抛异常。
实现自定义 Serdes
如果内置 Serde 无法满足需求,最佳起点是研读现有 Serde 的源码(见上一节)。典型工作流如下:
- 为你的数据类型
T实现序列化器:实现 org.apache.kafka.common.serialization.Serializer 接口; - 为
T实现反序列化器:实现 org.apache.kafka.common.serialization.Deserializer 接口; - 为
T实现Serde:实现 org.apache.kafka.common.serialization.Serde 接口——可以手动实现(参考内置 Serde),也可以借助 Serdes 中的Serdes.serdeFrom(Serializer<T>, Deserializer<T>)便捷方法。
Serde接口(Serde.java)本身非常轻量:继承Closeable,包含默认空实现的configure(Map, boolean)与close(),以及必须实现的serializer()与deserializer()两个抽象方法。其 Javadoc 明确指出:实现该接口的类需要有无参构造函数。
关键限制(务必注意):
- 若想把自定义 Serde 用于
KafkaStreams的配置(即DEFAULT_KEY_SERDE_CLASS_CONFIG/DEFAULT_VALUE_SERDE_CLASS_CONFIG),必须实现为无泛型的具体类——因为配置是按类名字符串反射实例化的; - 若你的 Serde 类带泛型,或使用
Serdes.serdeFrom(Serializer<T>, Deserializer<T>)组合而成,则只能通过方法调用方式传入,例如builder.stream("topicName", Consumed.with(...))。
Kafka Streams DSL for Scala 的隐式 Serdes
在使用 Kafka Streams DSL for Scala 时,不需要也不支持配置默认 Serde。Serde 改由 Scala 隐式机制提供——官方为常见基本数据类型提供了默认的隐式 Serde 实现。详见 DSL API 文档中的 Implicit Serdes 与 User-Defined Serdes 章节。这显著减少了 Scala 用户的样板代码:常见类型开箱即用,特殊类型则通过自定义隐式值无缝接入。
小结
Serde 是 Kafka Streams 应用与 Kafka 字节世界之间的桥梁。掌握三个层次即可覆盖绝大多数场景:
- 默认配置:通过
StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG/DEFAULT_VALUE_SERDE_CLASS_CONFIG设置全局限定 Serde; - 显式覆盖:在
Consumed、Produced等 DSL 伴随类中按需传入,支持细粒度选择性覆盖; - 窗口与自定义:窗口场景使用
WindowedSerdes,业务场景可参照内置实现编写自定义 Serde,并通过Serdes.serdeFrom(...)快速组合。
无论选择哪种方式,都请记住底层原则:配置路径需要无参构造的具体类,方法传参路径则更灵活。理解了这一点,你就能从容处理 Kafka Streams 中的一切序列化需求。
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考