【问题标题】:Analysis of Log with Spark Streaming使用 Spark Streaming 分析日志
【发布时间】:2018-11-25 07:43:03
【问题描述】:

我最近使用 Spark SQL 对静态日志文件进行了分析(找出出现十多次的 ip 地址等内容)。问题是from this site。但我使用了我自己的实现。我将日志读入 RDD,将该 RDD 转换为 DataFrame(在 POJO 的帮助下)并使用 DataFrame 操作。

现在我应该使用 Spark Streaming 对 30 分钟的流式日志文件以及一天的汇总结果进行类似的分析。可以再次找到解决方案here,但我想以另一种方式进行。所以我做的是这个

  • 使用 Flume 将数据从日志文件写入 HDFS 目录

  • 使用 JavaDStream 从 HDFS 读取 .txt 文件

然后我不知道如何进行。这是我使用的代码

Long slide = 10000L; //new batch every 10 seconds
Long window = 1800000L; //30 mins
SparkConf conf = new SparkConf().setAppName("StreamLogAnalyzer");
JavaStreamingContext streamingContext = new JavaStreamingContext(conf, new Duration(slide));
JavaDStream<String> dStream = streamingContext.textFileStream(hdfsPath).window(new Duration(window), new Duration(slide));

现在我似乎无法决定是否应该将每个批次转换为 DataFrame 并执行我之前对静态日志文件所做的操作。还是这种方式既费时又费力。

我是流媒体和 Flume 的绝对菜鸟。有人可以指导我吗?

【问题讨论】:

    标签: hdfs spark-streaming flume


    【解决方案1】:

    在 Spark 中使用 DataFrame(和 Dataset)是最新版本的 Spark 中最受推崇的方式,因此它是一个正确的选择。我认为当您将文件移动到 HDFS 而不是从任何事件日志中读取时,由于流的非显式性质,会出现一些模糊性。

    这里的重点是选择正确的批处理时间大小(或在您的 sn-p 中的幻灯片大小),因此应用程序将处理它在该时间槽下加载的数据并且不会有批处理队列。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-05-13
      • 1970-01-01
      • 1970-01-01
      • 2010-11-25
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多