ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

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

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提供 SerdeSerializer/Deserializer序列化器/反序列化器。本文以官方开发者指南 datatypes.md 为骨架结合 Kafka 仓库源码系统讲解 Serde 的配置方式、内置 Serde 清单、窗口 SerdeWindowed Serdes、自定义 Serde 实现以及 Scala DSL 的隐式 Serde 机制。读完后你将能熟练地在自己的 Streams 应用中配置、覆盖、组合与定制 Serde并理解其底层实现原理。说明docs/documentation/streams/developer-guide/datatypes.md为文档站点的重定向占位页真实内容位于 docs/streams/developer-guide/datatypes.md本文以其真实内容为准。为什么每个 Kafka Streams 应用都必须提供 SerdesKafka 本身只关心字节数组Producer 把任意类型序列化成byte[]写入 TopicConsumer 再把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 SerdeString stringSerde Serdes.String(); final SerdeLong longSerde Serdes.Long(); // KStream userCountByRegion 的键是 String区域值是 Long用户数 KStreamString, 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 SerdeLong longSerde Serdes.Long(); KStreamString, 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 groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version4.3.0/version /dependency这些实现位于org.apache.kafka.common.serialization包源码见 clients/src/main/java/org/apache/kafka/common/serialization可通过Serdes工厂类获取数据类型Serdebyte[]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())。WrapperSerde是Serde接口的一个便捷实现把configure/close/serializer/deserializer四个方法全部委托给内部的 Serializer 与 Deserializer这也解释了为什么大多数内置 Serde 本身不包含业务逻辑——它们只是“序列化器 反序列化器”的组合器。此外Serdes.serdeFrom(ClassT type)可以根据类型自动推断内置 Serde支持 String、Short、Integer、Long、Float、Double、byte[]、ByteBuffer、Bytes、UUID、Boolean遇到未知类型会抛出IllegalArgumentException而Serdes.serdeFrom(SerializerT, DeserializerT)则可以从独立的序列化器与反序列化器组合出任意 Serde。ListSerde值得一提它有两种构造方式——无参Serdes.ListSerde()直接使用内置的ListSerializer/ListDeserializer以及带参数Serdes.ListSerde(ClassL listClass, SerdeInner innerSerde)后者允许指定 List 的具体实现类如ArrayList.class与元素类型的内置 Serde从而支持泛型列表的序列化。JSONKafka Streams 官方代码示例中包含一个基础的 JSON Serde 实现PageViewTypedDemo.java如示例所示可以通过Serdes.serdeFrom(serializer实例, deserializer实例)组合出 JSON 兼容的序列化器与反序列化器。示例中的JSONSerdeT类同时实现了SerializerT、DeserializerT与SerdeT三个接口内部基于 Jackson 的ObjectMapper反序列化时用JsonTypeInfo注解记录的_t字段区分具体子类型配合JsonSubTypes注册所有可能的类型。这种“泛型 JSON Serde 类型字段”的模式非常适合一类数据对应多种具体结构的场景。窗口 SerdesWindowed Serdes时间窗口Tumbling、Hopping、Session Window是 Kafka Streams 最核心的抽象之一。窗口操作如windowedBy、groupByKey 窗口聚合产生的键类型是WindowedK——它把“原始键 窗口起始/结束时间”打包在一起。kafka-streams构件为这类窗口类型提供了专门的 Serde 实现dependency groupIdorg.apache.kafka/groupId artifactIdkafka-streams/artifactId version4.3.0/version /dependency这些实现位于org.apache.kafka.streams.kstream包源码见 streams/src/main/java/org/apache/kafka/streams/kstreamSerdes组合类WindowedSerdes.TimeWindowedSerdeTWindowedSerdes.SessionWindowedSerdeT序列化器TimeWindowedSerializerTSessionWindowedSerializerT反序列化器TimeWindowedDeserializerTSessionWindowedDeserializerT从源码看TimeWindowedSerializer.serialize()最终委托给 WindowKeySchema.toBinary(...) 完成字节编码TimeWindowedDeserializer.deserialize()则根据是否 changelog Topic 分别调用WindowKeySchema.fromStoreKey(...)或WindowKeySchema.from(...)。因此窗口 Serde 的底层字节布局与状态存储State Store的键模式完全一致这也保证了窗口聚合结果能无缝落盘与恢复。代码中的用法// 时间窗口 Serde —— 工厂方法 SerdeWindowedString timeWindowedSerde WindowedSerdes.timeWindowedSerdeFrom(String.class, 500L); // 时间窗口 Serde —— 构造函数 SerdeWindowedString timeWindowedSerde2 new WindowedSerdes.TimeWindowedSerde(Serdes.String(), 500L); // 会话窗口 Serde —— 工厂方法 SerdeWindowedString sessionWindowedSerde WindowedSerdes.sessionWindowedSerdeFrom(String.class); // 会话窗口 Serde —— 构造函数 SerdeWindowedString sessionWindowedSerde2 new WindowedSerdes.SessionWindowedSerde(Serdes.String()); // 单独使用序列化器 / 反序列化器 TimeWindowedSerializerString serializer new TimeWindowedSerializer(Serdes.String().serializer()); TimeWindowedDeserializerString deserializer new TimeWindowedDeserializer(Serdes.String().deserializer(), 500L);窗口大小如上面示例中的500L单位毫秒对时间窗口 Serde 至关重要反序列化器需要它来计算窗口的结束时间窗口的字节编码中通常只保存起始时间戳。会话窗口没有固定大小因此不需要该参数。命令行工具中的用法使用命令行工具如bin/kafka-console-consumer.sh消费窗口结果 Topic 时可以通过--formatter-property传入窗口反序列化器与窗口大小属性名遵循前缀模式# 时间窗口反序列化器配置 --formatter-property print.keytrue \ --formatter-property key.deserializerorg.apache.kafka.streams.kstream.TimeWindowedDeserializer \ --formatter-property key.deserializer.windowed.inner.deserializer.classorg.apache.kafka.common.serialization.StringDeserializer \ --formatter-property key.deserializer.window.size.ms500 # 会话窗口反序列化器配置 --formatter-property print.keytrue \ --formatter-property key.deserializerorg.apache.kafka.streams.kstream.SessionWindowedDeserializer \ --formatter-property key.deserializer.windowed.inner.deserializer.classorg.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.classStreamsConfig.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”必须一致否则抛出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(SerializerT, DeserializerT)便捷方法。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(SerializerT, DeserializerT)组合而成则只能通过方法调用方式传入例如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),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进