
在使用 TopologyTestDriver 进行 Kafka Streams 单元测试时,必须显式配置 DEFAULT_KEY_SERDE_CLASS_CONFIG 和 DEFAULT_VALUE_SERDE_CLASS_CONFIG,否则即使为具体类型指定了自定义 Serde,也会因缺少全局默认序列化器而抛出 ConfigException。
在使用 topologytestdriver 进行 kafka streams 单元测试时,必须显式配置 `default_key_serde_class_config` 和 `default_value_serde_class_config`,否则即使为具体类型指定了自定义 serde,也会因缺少全局默认序列化器而抛出 configexception。
Kafka Streams 在构建拓扑(尤其是涉及状态存储、窗口操作或内部处理器初始化)时,会依赖 StreamsConfig 中的默认键/值 Serde 作为兜底策略。即使你在 createInputTopic() 中为每个 Topic 显式传入了 serializer()(如 myObject1Serde.serializer()),Kafka Streams 内部仍会在初始化状态存储(如 MeteredWindowStore)、处理器上下文(AbstractProcessorContext)等环节尝试通过 StreamsConfig#defaultValueSerde() 获取默认值 Serde —— 若该配置未设置,就会触发你遇到的 ConfigException: Please specify a value serde or set one through StreamsConfig#DEFAULT_VALUE_SERDE_CLASS_CONFIG。
关键点在于:createInputTopic() 的参数仅影响测试驱动的数据注入行为,不参与 Streams 运行时的全局 Serde 解析逻辑。因此,必须在 TopologyTestDriver 初始化前,将默认 Serde 配置注入 Properties 对象,并传递给拓扑构建器(或直接用于 TopologyTestDriver 构造)。
✅ 正确做法如下(适配你的测试场景):
@BeforeEach
void setUp() {
// 1. 构建必需的 Streams 配置(即使本地测试也需完整)
final Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "test-app");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "dummy:9092"); // 测试用,可为任意字符串
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());
// 2. 使用配置构建 TopologyTestDriver(注意:需传入 props)
testDriver = new TopologyTestDriver(myTopology.buildTopology(), props);
// 3. 定义具体类型的 Serde(保持不变)
Serde<myobject1> myObject1Serde = new JsonSerde(MyObject1.class);
Serde<myobject2> myObject2Serde = new JsonSerde(MyObject2.class);
Serde<myobject3> myObject3Serde = new JsonSerde(MyObject3.class);
// 4. 创建输入/输出 Topic(此时 serializer 已生效,且全局默认 Serde 已就绪)
myObject1Topic = testDriver.createInputTopic(
MyTopology.FIRST_INPUT_TOPIC,
Serdes.String().serializer(),
myObject1Serde.serializer()
);
// 同理创建其他 topic...
}</myobject3></myobject2></myobject1>
⚠️ 注意事项:
-
DEFAULT_VALUE_SERDE_CLASS_CONFIG必须是 Serde 实现类的Class对象(如Serdes.String().getClass()),而非实例(如new JsonSerde(String.class)),否则会因类型不匹配导致运行时异常。 - 若拓扑中使用了泛型类型(如
KTable<string myobject1></string>),且未显式指定Materialized.with(...),Streams 会回退到默认 Serde,因此建议对关键状态存储显式声明 Serde:table.toStream().to("output-topic", Produced.with(Serdes.String(), myObject3Serde)); -
BOOTSTRAP_SERVERS_CONFIG在TopologyTestDriver中虽不真正连接集群,但必须提供非空值(Kafka 3.0+ 仍校验该属性),推荐使用"dummy:9092"或"localhost:1"等占位符。
? 总结:Kafka Streams 的 Serde 配置分两级——全局默认(必配) 与 局部显式(按需)。单元测试不是“免配置环境”,而是更需严谨模拟生产配置。只要补全 StreamsConfig 中的两个默认 Serde 类配置,你的合并拓扑即可顺利初始化并完成端到端验证。











