You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于文件内容在FileStream Source Connector中路由到不同Topic

问题分析与解决

你的FileStream Source Connector配置中,RegexRouter转换的分组引用错误,导致无法根据内容中的关键字路由到对应Topic。

原配置问题点

原配置里的transforms.routeToTopic.replacement=$1引用的是第一个捕获组(.*)(即关键字前面的所有内容),而非我们需要提取的info/warn/error关键字,所以所有消息都默认 fallback 到了default-topic。

修正后的完整配置

name=my-filestream-connector
connector.class=FileStreamSource
tasks.max=1
file=D:/kafka/kafka/config/demo_file4.txt
value.converter=org.apache.kafka.connect.storage.StringConverter
key.converter=org.apache.kafka.connect.storage.StringConverter
topic=default-topic
transforms=routeToTopic
transforms.routeToTopic.type=org.apache.kafka.connect.transforms.RegexRouter
transforms.routeToTopic.regex=(.*)(info|warn|error)(.*)
transforms.routeToTopic.replacement=$2
transforms.routeToTopic.topic.format=%s

关键修改说明

  • 将transforms.routeToTopic.replacement的值从$1改为$2:$2对应正则中的第二个捕获组(info|warn|error),也就是要提取的关键字,消息会被路由到以该关键字命名的Topic(如info、warn、error)。
  • 保留topic=default-topic:当消息内容不包含这三个关键字时,会默认发送到这个Topic,你可根据需求调整或移除。

验证效果

向文件写入以下内容时:

2024-05-20 14:30:00 INFO: User login success
2024-05-20 14:31:00 ERROR: Database connection failed
2024-05-20 14:32:00 WARN: Disk usage is high
2024-05-20 14:33:00 Debug: System status check

前三条消息会分别发送到info、error、warn Topic,最后一条不包含关键字的消息会发送到default-topic。

内容的提问来源于stack exchange,提问作者Raghav07

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 21:37:16