如何在Scala中持续读取文件新增内容并推送至Kafka生产者?
实现Kafka监控文件新增内容并推送的正确姿势
嘿,我注意到你当前的代码有个关键问题:你虽然执行了tail -f命令,但完全没用到这个进程的输出,反而用Source.fromFile去读取整个文件的内容——这样不仅会把文件里已有的所有内容都发送一遍,而且根本不会监听后续新增的内容!
下面给你修正后的实现方案,核心是读取tail -f进程的标准输出流,这样就能持续获取文件的新增内容:
import org.apache.kafka.clients.producer.{KafkaProducer, ProducerRecord} import java.io.BufferedReader import java.io.InputStreamReader import java.util.Properties object FileMonitorKafkaProducer { def main(args: Array[String]): Unit = { val fileName = "abc.txt" val topicName = "topicName" // 配置Kafka生产者(替换成你的实际集群地址) val props = new Properties() props.put("bootstrap.servers", "localhost:9092") props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer") props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer") val producer = new KafkaProducer[String, String](props) var process: Process = null var reader: BufferedReader = null try { // 用数组形式执行命令,避免文件名含特殊字符时出错 val command = Array("tail", "-f", fileName) process = Runtime.getRuntime.exec(command) // 读取tail命令的实时输出流 reader = new BufferedReader(new InputStreamReader(process.getInputStream)) var line: String = null // 持续读取新增的每一行内容 while ({line = reader.readLine(); line != null}) { val message = line + "\n" val producerRecord = new ProducerRecord[String, String](topicName, message) producer.send(producerRecord) println(s"已推送内容到Kafka: $line") } } catch { case e: Exception => e.printStackTrace() } finally { // 确保所有资源被正确关闭,避免泄漏 if (reader != null) reader.close() if (process != null) process.destroy() producer.close() } } }
几个关键细节说明:
- 把
tail -f命令改成数组形式,能避免文件名包含空格、特殊字符时出现解析错误。 - 必须读取
Process的getInputStream(),这才是tail -f输出新增内容的通道,之前的代码完全忽略了这一点。 - 在
finally块统一关闭资源,防止进程残留或内存泄漏。
可选的跨平台替代方案:
如果你的程序需要在Windows等非类Unix环境运行,依赖tail命令就不太友好了。这时候可以用Java的WatchService监听文件所在目录的修改事件,记录上次读取的文件位置,每次触发事件时从该位置读取新增内容,再推送到Kafka。不过如果你的场景只在类Unix环境下运行,tail -f的方案足够简单高效。
内容的提问来源于stack exchange,提问作者Naman Agarwal
相关产品推荐
相关产品推荐

