Advertisement

学习spark | spark streaming和kafka之间结合进行实时的数据处理(基于文件流)

阅读量:

本篇文章分两个部分:

一个为已完成的生产者与消费者代码示例,另一个则为对代码实现流程的说明

1. 首先完善好的代码

生产者代码

复制代码
    import java.io.{File, RandomAccessFile}
    import java.nio.charset.StandardCharsets
    import scala.io.Source
    
    object KafkaWordProducer2 {
      def main(args: Array[String]) {
    if (args.length < 3) {
      System.err.println("用法: KafkaWordProducer <metadataBrokerList> <topic> <linesPerSec>")
      System.exit(1)
    }
    
    val Array(brokers, topic, linesPerSec) = args
    // Kafka生产者属性
    val props = new java.util.HashMap[String, Object]()
    props.put(org.apache.kafka.clients.produ

全部评论 (0)

还没有任何评论哟~