news 2026/9/10 14:59:10

Apache Kafka Streams 数据类型与序列化(Serdes)完全指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Kafka Streams 数据类型与序列化(Serdes)完全指南

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 有且仅有两种方式(至少使用其中一种):

  1. java.util.Properties配置实例中设置默认 Serde(通过StreamsConfig的两个配置项);
  2. 在调用相应 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 参数的是ConsumedProducedGroupedJoinedMaterialized等“伴随类”,它们都有静态工厂方法与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()(见下方提示)
ByteBufferSerdes.ByteBuffer()
DoubleSerdes.Double()
IntegerSerdes.Integer()
LongSerdes.Long()
StringSerdes.String()
UUIDSerdes.UUID()
VoidSerdes.Void()
ListSerdes.ListSerde()
BooleanSerdes.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())WrapperSerdeSerde接口的一个便捷实现,把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 最核心的抽象之一。窗口操作(如windowedBygroupByKey+ 窗口聚合)产生的键类型是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/SessionWindowedDeserializerwindowed.inner.deserializer.class指定窗口内部键的反序列化器,window.size.ms指定窗口大小(毫秒)。会话窗口无需window.size.ms

已弃用的配置

以下StreamsConfig参数已弃用,应改用向序列化器/反序列化器构造函数直接传参的方式:

  • StreamsConfig.WINDOWED_INNER_CLASS_SERDE→ 改用TimeWindowedSerializer.WINDOWED_INNER_SERIALIZER_CLASSwindowed.inner.serializer.class)与TimeWindowedDeserializer.WINDOWED_INNER_DESERIALIZER_CLASSwindowed.inner.deserializer.class
  • StreamsConfig.WINDOW_SIZE_MS_CONFIG→ 改用TimeWindowedDeserializer.WINDOW_SIZE_MS_CONFIGwindow.size.ms

从 TimeWindowedSerializer.java 的configure()实现可以看出兼容逻辑:它优先读取windowed.inner.serializer.class,若为空则回退读取已弃用的StreamsConfig.WINDOWED_INNER_CLASS_SERDE并打印弃用告警;同时校验“构造函数传入的 inner 与配置指定的 inner”必须一致,否则抛出IllegalArgumentExceptionTimeWindowedDeserializer的窗口大小解析也有类似的约束:构造函数与window.size.ms配置只能二选一,二者都设置或都不设置都会抛异常

实现自定义 Serdes

如果内置 Serde 无法满足需求,最佳起点是研读现有 Serde 的源码(见上一节)。典型工作流如下:

  1. 为你的数据类型T实现序列化器:实现 org.apache.kafka.common.serialization.Serializer 接口;
  2. T实现反序列化器:实现 org.apache.kafka.common.serialization.Deserializer 接口;
  3. 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 字节世界之间的桥梁。掌握三个层次即可覆盖绝大多数场景:

  1. 默认配置:通过StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG/DEFAULT_VALUE_SERDE_CLASS_CONFIG设置全局限定 Serde;
  2. 显式覆盖:在ConsumedProduced等 DSL 伴随类中按需传入,支持细粒度选择性覆盖;
  3. 窗口与自定义:窗口场景使用WindowedSerdes,业务场景可参照内置实现编写自定义 Serde,并通过Serdes.serdeFrom(...)快速组合。

无论选择哪种方式,都请记住底层原则:配置路径需要无参构造的具体类,方法传参路径则更灵活。理解了这一点,你就能从容处理 Kafka Streams 中的一切序列化需求。

【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/10 14:58:43

COMSOL中手性介质的电磁仿真与应用

1. 手性介质在电磁仿真中的独特价值手性介质&#xff08;Chiral media&#xff09;是一类具有特殊电磁响应的材料&#xff0c;其本构关系中电场与磁场存在交叉耦合。这种特性使得电磁波在传播时会发生偏振面旋转&#xff0c;这种现象被称为光学活性。在COMSOL Multiphysics中模…

作者头像 李华
网站建设 2026/9/10 14:58:39

GTD时间管理法:提升个人生产力的核心技巧

1. 项目概述&#xff1a;为什么《尽管去做》值得一读&#xff1f;这本书的核心价值在于它提供了一套完整的个人生产力管理系统&#xff0c;帮助读者从"想法积压"的状态转变为"高效执行"的模式。作者David Allen提出的GTD&#xff08;Getting Things Done&a…

作者头像 李华
网站建设 2026/9/10 14:54:58

Elasticsearch核心概念与生产环境部署指南

1. Elasticsearch核心概念全景解析 Elasticsearch作为当前最流行的分布式搜索和分析引擎&#xff0c;其核心架构设计与传统数据库有着本质区别。我在实际项目中使用ES近五年&#xff0c;发现许多开发者最初接触时容易被其术语体系迷惑。这里我将用最直白的语言拆解这些概念。 …

作者头像 李华
网站建设 2026/9/10 14:53:40

笔墨AI真的能搞定毕业论文?实测真相来了

毕业季论文工具层出不穷&#xff0c;但大多功能单一、质量参差不齐&#xff0c;要么生成内容空洞&#xff0c;要么查重率超标&#xff0c;很难满足高校论文考核标准。在众多学术AI工具中&#xff0c;笔墨AI凭借垂直化、专业化的学术服务脱颖而出&#xff0c;成为大四学生的热门…

作者头像 李华