【发布时间】: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