Kafka Connect FileStreamSource忽略新增行问题求助
看起来你遇到的核心问题是偏移量存储导致的读取进度异常,结合你的操作步骤和日志来看,我来一步步帮你解决:
问题根源分析
你使用的是独立模式的Kafka Connect,偏移量被存储在/tmp/connect.offsets文件中。FileStreamSource连接器是通过记录文件的字节偏移量来跟踪读取进度的,当你做了以下操作后:
- 删除了原主题并创建新主题
- 修改了日志文件的内容(删除部分行)
- 更换了连接器名称但未重置偏移量
偏移量文件中可能还保留着针对原日志文件的旧偏移记录——这个偏移量可能已经大于当前日志文件的实际大小。当新连接器启动时,它读取到这个旧偏移量,会认为已经读到了文件末尾,因此完全忽略后续新增的内容(这也对应了你日志中反复出现的flushing 0 outstanding messages)。
具体解决步骤
1. 停止Kafka Connect服务
先确保Connect进程完全停止,避免修改偏移量文件时出现冲突:
# 如果是前台启动的直接按Ctrl+C,否则用进程杀死命令 pkill -f connect-standalone.sh
2. 重置偏移量
因为你用的是文件存储偏移量,最简单的方式就是删除这个偏移量文件:
rm /tmp/connect.offsets
如果你不想删除整个文件,也可以手动编辑它(它是JSON格式),找到对应new-file-connector或者目标日志文件的偏移记录,把offset字段的值改为0,这样连接器会从头开始读取文件,并且会正常跟踪后续新增内容。
3. 确认日志文件的一致性
确保你修改的日志文件是同一个物理文件(没有被删除后重新创建),可以通过查看inode来验证:
ls -i /data/users/zamara/suivi_prod/app/data/logs.txt
如果文件被删除重建,inode会发生变化,FileStreamSource可能无法识别为同一个文件,这种情况也会导致读取异常,需要确保是原文件或者重新指定配置中的文件路径。
4. 重新启动Kafka Connect
使用你的新配置重新启动Connect:
connect-standalone.sh config/connect-standalone.properties config/connect-file-source.properties
5. 验证结果
现在往日志文件里新增一行内容,然后用控制台消费者查看:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic new-topic --from-beginning
你应该能看到新增的行被正常消费到了。
生产环境建议
- 尽量使用分布式模式的Kafka Connect,偏移量会存储在Kafka内部主题中,管理更灵活,也避免了文件存储的单点问题。
- FileStreamSource是官方提供的演示用连接器,不适合生产环境。生产环境可以考虑使用更可靠的文件源连接器,或者自定义连接器来处理日志的轮转、内容变更等场景。
- 避免直接修改正在被读取的日志文件,建议用
logrotate这类工具做日志轮转,让连接器能更好地跟踪文件变化。
内容的提问来源于stack exchange,提问作者Alexis

