使用CloudEvent实现KStream与KTable关联失败求助
我在Kafka Streams中实现KStream与KTable的关联操作,使用普通流数据时POC运行正常,但使用CloudEvent时始终遇到序列化相关问题。
代码示例
Map<String, Object> ceSerializerConfigs = new HashMap<>(); ceSerializerConfigs.put(ENCODING_CONFIG, Encoding.STRUCTURED); ceSerializerConfigs.put(EVENT_FORMAT_CONFIG, JsonFormat.CONTENT_TYPE); CloudEventSerializer serializer = new CloudEventSerializer(); serializer.configure(ceSerializerConfigs, false); CloudEventDeserializer deserializer = new CloudEventDeserializer(); deserializer.configure(ceSerializerConfigs, false); Serde<CloudEvent> cloudEventSerde = Serdes.serdeFrom(serializer, deserializer); KStream<String, CloudEvent> kStream = builder.stream("stream-topic", Consumed.with(Serdes.String(), cloudEventSerde)); KTable<String, CloudEvent> kTable = builder.table("ktable-topic", Consumed.with(Serdes.String(), cloudEventSerde)); KStream<String, CloudEvent> joined = kStream .join(kTable, (left, right) -> CloudEventBuilder.v1().withId(left.getId().concat(right.getId())).build()); joined.to(output, Produced.with(Serdes.String(), eventsSerde)); KafkaStreams kafkaStreams = new KafkaStreams(builder.build(), streamProps); kafkaStreams.start();
尝试过的方案及异常
我尝试过使用WrapperSerde,但仍持续抛出以下异常:
18:12:08.691
[basic-streams-updated-0630c691-0080-4e02-8c85-7bff650f34e9-StreamThread-1]
ERROR org.apache.kafka.streams.KafkaStreams - stream-client
[basic-streams-updated-0630c691-0080-4e02-8c85-7bff650f34e9]
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: Exception caught in
process. taskId=0_0, processor=KSTREAM-SOURCE-0000000002,
topic=cloudevent-ktable, partition=0, offset=80,
stacktrace=java.lang.UnsupportedOperationException:
CloudEventSerializer supports only the signature serialize(String,
Headers, CloudEvent)Caused by: java.lang.UnsupportedOperationException:
CloudEventSerializer supports only the signature serialize(String,
Headers, CloudEvent) at
io.cloudevents.kafka.CloudEventSerializer.serialize(CloudEventSerializer.java:84)
~[cloudevents-kafka-2.5.0.jar:?] at
io.cloudevents.kafka.CloudEventSerializer.serialize(CloudEventSerializer.java:38)
~[cloudevents-kafka-2.5.0.jar:?] at
org.apache.kafka.streams.state.internals.ValueAndTimestampSerializer.serialize(ValueAndTimestampSerializer.java:82)
~[kafka-streams-2.8.0.jar:?] at
org.apache.kafka.streams.state.internals.ValueAndTimestampSerializer.serialize(ValueAndTimestampSerializer.java:73)
~[kafka-streams-2.8.0.jar:?] at
org.apache.kafka.streams.state.internals.ValueAndTimestampSerializer.serialize(ValueAndTimestampSerializer.java:30)
~[kafka-streams-2.8.0.jar:?] at
org.apache.kafka.streams.state.StateSerdes.rawValue(StateSerdes.java:192)
~[kafka-streams-2.8.0.jar:?] at
org.apache.kafka.streams.state.internals.MeteredKeyValueStore.lambda$put$4(MeteredKeyValueStore.java:200)
~[kafka-streams-2.8.0.jar:?] at
org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:884)
~[kafka-streams-2.8.0.jar:?] at
org.apache.kafka.streams.state.internals.MeteredKeyValueStore.put(MeteredKeyValueStore.java:200)
~[kafka-streams-2.8.0.jar:?] at
org.apache.kafka.streams.processor.internals.AbstractReadWriteDecorator$KeyValueStoreReadWriteDecorator.put(AbstractReadWriteDecorator.java:120)
~[kafka-streams-2.8.0.jar:?] at
org.apache.kafka.streams.kstream.internals.KTableSource$KTableSourceProcessor.process(KTableSource.java:122)
~[kafka-streams-2.8.0.jar:?]
请问有人成功将CloudEvent与KTable配合使用过吗?
异常原因分析
原生CloudEventSerializer仅支持带Headers参数的serialize(String, Headers, CloudEvent)签名,但Kafka Streams处理KTable状态存储时,会通过ValueAndTimestampSerializer调用无Headers的serialize(String, T)方法,直接触发了序列化器的不支持异常。
具体解决步骤
1. 自定义兼容的CloudEvent序列化/反序列化器
包装原生CloudEvent序列化器,补全无Headers的方法实现:
// 自定义序列化器 public class CloudEventKafkaSerializer implements Serializer<CloudEvent> { private final CloudEventSerializer delegate = new CloudEventSerializer(); @Override public void configure(Map<String, ?> configs, boolean isKey) { delegate.configure(configs, isKey); } @Override public byte[] serialize(String topic, CloudEvent data) { // 调用带Headers的重载,传入空Headers对象 return delegate.serialize(topic, new RecordHeaders(), data); } @Override public byte[] serialize(String topic, Headers headers, CloudEvent data) { return delegate.serialize(topic, headers, data); } @Override public void close() { delegate.close(); } } // 自定义反序列化器(可选,确保兼容无Headers场景) public class CloudEventKafkaDeserializer implements Deserializer<CloudEvent> { private final CloudEventDeserializer delegate = new CloudEventDeserializer(); @Override public void configure(Map<String, ?> configs, boolean isKey) { delegate.configure(configs, isKey); } @Override public CloudEvent deserialize(String topic, byte[] data) { return delegate.deserialize(topic, new RecordHeaders(), data); } @Override public CloudEvent deserialize(String topic, Headers headers, byte[] data) { return delegate.deserialize(topic, headers, data); } @Override public void close() { delegate.close(); } }
2. 创建兼容的Serde
使用自定义的序列化器创建CloudEvent Serde:
Map<String, Object> ceSerializerConfigs = new HashMap<>(); ceSerializerConfigs.put(ENCODING_CONFIG, Encoding.STRUCTURED); ceSerializerConfigs.put(EVENT_FORMAT_CONFIG, JsonFormat.CONTENT_TYPE); CloudEventKafkaSerializer serializer = new CloudEventKafkaSerializer(); serializer.configure(ceSerializerConfigs, false); CloudEventKafkaDeserializer deserializer = new CloudEventKafkaDeserializer(); deserializer.configure(ceSerializerConfigs, false); Serde<CloudEvent> cloudEventSerde = Serdes.serdeFrom(serializer, deserializer);
3. 显式指定KTable状态存储的Serde
创建KTable时,通过Materialized参数显式指定状态存储使用的Serde,避免默认适配逻辑:
KTable<String, CloudEvent> kTable = builder.table( "ktable-topic", Consumed.with(Serdes.String(), cloudEventSerde), Materialized.<String, CloudEvent, KeyValueStore<Bytes, byte[]>>as("cloud-event-ktable-store") .withKeySerde(Serdes.String()) .withValueSerde(cloudEventSerde) );
4. 统一输出Serde
确保关联后的输出流使用同一个兼容的Serde(修正原代码中可能的笔误):
joined.to(output, Produced.with(Serdes.String(), cloudEventSerde));
原理说明
自定义序列化器通过补全无Headers的方法实现,将调用转发到原生支持的带Headers重载;同时显式指定状态存储Serde,让Kafka Streams在处理KTable状态时直接使用我们的兼容序列化器,绕过了默认的ValueAndTimestampSerializer适配逻辑,从而解决序列化异常。
内容的提问来源于stack exchange,提问作者SWA PSTG

