如何在KafkaConnector中按Avro Key字段过滤消息?现有方案无效
问题描述
我有一个将Kafka记录下沉至PubSub主题的KafkaConnector,其Key为Avro格式,Value为字节类型。希望添加转换器,基于Key中的first字段过滤消息。
现有KafkaConnector的YAML配置:
apiVersion: kafka.strimzi.io/v1beta2 kind: KafkaConnector metadata: name: my-kafka-connector labels: strimzi.io/cluster: connect-pubsub-sink spec: class: com.google.pubsub.kafka.sink.CloudPubSubSinkConnector tasksMax: 2 config: schema.registry.url: "${env:SCHEMA_REGISTRY_URL}" key.converter: "io.confluent.connect.avro.AvroConverter" key.converter.schema.registry.url: "${env:SCHEMA_REGISTRY_URL}" value.converter: "org.apache.kafka.connect.converters.ByteArrayConverter" value.converter.schema.registry.url: "${env:SCHEMA_REGISTRY_URL}" cps.topic: "output_topic" cps.project: "${env:GCP_PROJECT}" orderingKeySource: "key" headers.publish: "true" topics: "input-topic"
Key的Avro Schema字段:
first String, second String, third String, fourth String
尝试过用TopicNameMatches断言结合自定义转换器的方案但无效,配置如下:
transforms: "Concatenate,Extract" transforms.Concatenate.type: "phenix.connect.transformers.ConcatenateFields$Key" transforms.Concatenate.field.list: "first,second,third,fourth" transforms.Concatenate.field.new: "concatenatedKey" transforms.Concatenate.delimiter: "#" transforms.Extract.type: "org.apache.kafka.connect.transforms.ExtractField$Key" transforms.Extract.field: "concatenatedKey" transforms.Extract.predicate: "Filter" predicates: "Filter" predicates.Filter.type: "org.apache.kafka.connect.transforms.predicates.TopicNameMatches" predicates.Filter.pattern: "^FIRST-00#.+$"
请问是否有可行的实现方式?还是必须开发Kafka Streams应用才能实现这一基础功能?
可行方案
不需要开发Kafka Streams应用,有几种直接的方式实现基于Avro Key字段的过滤:
方案1:使用Confluent的FieldValueMatches断言(推荐)
Confluent Platform提供的FieldValueMatches断言可直接匹配记录Key/Value中的字段值,完全适配你的需求。
配置示例
在Connector的config节点中添加以下配置:
# 定义过滤转换器 transforms: FilterRecords transforms.FilterRecords.type: org.apache.kafka.connect.transforms.Filter$Key transforms.FilterRecords.predicate: MatchFirstField # 定义断言规则 predicates: MatchFirstField predicates.MatchFirstField.type: org.apache.kafka.connect.transforms.predicates.FieldValueMatches$Key # 指定要匹配的Key字段 predicates.MatchFirstField.field: first # 设置匹配正则(示例:仅保留first字段值为"FIRST-00"的记录) predicates.MatchFirstField.pattern: ^FIRST-00$ # 可选:设置为true时会过滤掉匹配的记录(保留不匹配的),默认false(保留匹配的) # predicates.MatchFirstField.negate: false
方案2:使用开源第三方断言插件
如果无法使用Confluent商业组件,可选择开源的第三方断言插件(如kafka-connect-field-value-filter),或自行编写简易自定义Predicate:
- 实现
org.apache.kafka.connect.transforms.predicates.Predicate接口 - 在
test方法中解析Avro Key的first字段,判断是否符合过滤条件 - 将打包后的jar放到Connect插件目录,即可在配置中引用
方案3:原生Transform组合(无商业组件依赖)
仅用Apache Kafka原生Transforms的话,可先将Key的first字段提取到记录Header,再基于Header值过滤:
配置示例
transforms: ExtractFirstToHeader,FilterByHeader # 第一步:将Key的first字段提取到自定义Header transforms.ExtractFirstToHeader.type: org.apache.kafka.connect.transforms.InsertField$Key transforms.ExtractFirstToHeader.header.field: key_first transforms.ExtractFirstToHeader.static.field: ${kafka.key.first} # 第二步:基于Header值过滤记录 transforms.FilterByHeader.type: org.apache.kafka.connect.transforms.Filter$Value transforms.FilterByHeader.predicate: MatchHeaderFirst predicates: MatchHeaderFirst predicates.MatchHeaderFirst.type: org.apache.kafka.connect.transforms.predicates.HeaderValueMatches predicates.MatchHeaderFirst.header: key_first predicates.MatchHeaderFirst.pattern: ^FIRST-00$
注意:需验证Connect版本对该逻辑的兼容性。
内容的提问来源于stack exchange,提问作者Bilal Ennouali
相关产品推荐
相关产品推荐

