如何配置Kafka Source避免覆盖CloudEvent原始消息头
问题说明
目标是让携带原始Headers的CloudEvent格式Kafka消息,沿Kafka Source -> Broker -> ASP.NET Core服务链路完成传递,全程保留初始Kafka消息的所有Headers。
当前向Kafka投递包含消息体、自定义Headers的消息后,消息可被Kafka Source正常消费,但原始Kafka消息的Headers在Kafka到后端服务的传输链路中被替换,头信息对比如下:
初始投递的Kafka消息Headers
correlationid = {guid} ce-specversion = 1.0 ce-id = {guid} ce-source = {differentRelativeUriThanBelow} ce-type = {com.company.product.request.amqp.asynchronous:v1} Content-Type = application/cloudevents
服务端实际接收到的Headers
correlationid = ce-specversion = 1.0 ce-id = partition:0/offset:52 ce-source = /apis/v1/namespaces/myNamespace/kafkasources/kafka-source-myNamespace#myKafkaTopic ce-type = dev.knative.kafka.event Content-Type = application/cloudevents
需要确认是否存在方法阻止该默认覆盖行为,或通过合理配置让后端服务收到的HTTP请求中包含原始Kafka消息的所有Headers。
可行配置方案
Knative Kafka Source默认会将消费到的Kafka消息重新包装为Knative标准格式的CloudEvent,因此会默认覆盖ce-id、ce-source、ce-type这类标准CloudEvent字段,非ce-前缀的自定义头默认也不会主动透传,可通过以下配置实现原始头保留:
- 开启CloudEvent透传模式
给对应的Kafka Source资源添加注解kafka.eventing.knative.dev/passthrough: "true",开启该模式后Kafka Source不会对Kafka中已有的合法CloudEvent消息做重新包装,不会覆盖原始的标准CloudEvent字段,只会补充链路路由所需的必要元数据,原始消息携带的所有Headers都会沿链路透传。 - 配置指定头透传规则
如果不需要透传全量Headers,仅需保留部分自定义头,可以在Kafka Source的spec.ceOverrides.extensions字段中配置需要透传的Kafka消息头映射关系,配置后指定的头会作为CloudEvent扩展字段沿链路传递,不会出现值被置空的问题。
注意:透传模式开启后,链路中的Broker组件不会修改原始CloudEvent的标准字段,后端ASP.NET Core服务可直接从HTTP请求头中读取到所有原始Kafka消息携带的Headers。
内容的提问来源于stack exchange,提问作者Prox
相关产品推荐
相关产品推荐

