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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.20 10:20:26