【问题标题】:Persisting tweets using Spark Streaming使用 Spark Streaming 持久化推文
【发布时间】:2014-12-23 06:48:49
【问题描述】:

首先,我们的要求相当简单。当推文进来时,我们需要做的就是将它们保存在 HDFS 上(定期)。

JavaStreamingContext 的“检查点”API 看起来很有希望,但经过进一步审查,它似乎服务于不同的目的。 (另外,我不断收到 '/checkpoint/temp, error: No such file or directory (2)' 错误,但我们暂时不用担心)。

问题:JavaDStream 没有“saveAsHadoopFiles”方法——这有点道理。我想从流式作业保存到 Hadoop 不是一个好主意。

推荐的方法是什么?我是否应该将传入的“推文”写入 Kafka 队列,然后使用诸如“Camus”(https://github.com/linkedin/camus) 之类的工具推送到 HDFS?

【问题讨论】:

  • 为什么从 Streaming 作业保存到 hadoop 不是一个好主意?我认为这就是您真正想要的。
  • 每次收到消息,如果我们保存到 HDFS,我们的解决方案会扩展吗? Twitter 每秒发送数百万条推文。将每条推文直接插入 HDFS 不会扩展!会吗?
  • 如果 HDFS 的写入吞吐量无法保持持续的消息写入,那么在两者之间添加另一个系统(如 kafka)将如何提供帮助?使用经过调整的窗口(x 秒),您可以收集足够的消息以微批量写入 HDFS。这应该是相当有效的。
  • “调谐窗口”正是 Kafka 将提供给我们的,不是吗?此外,还有其他好处。 Storm 和 Spark 流都与 Kafka 很好地集成以进行实时处理。
  • Kafka 为您提供了一个高吞吐量的队列,但它是另一个增加系统复杂性的元素。如果您预期的瓶颈是 HDFS,我看不出 kafka 可以如何帮助您。

标签: twitter hdfs twitter4j apache-spark spark-streaming


【解决方案1】:

看到这个真棒博客文章证实了我的想法。作者使用 Kafka、Storm、Camus 等技术构建了一个“外汇交易系统”。这个用例与我的相似,所以我将使用这个设计和工具。谢谢。

http://insightdataengineering.com/blog/Building_a_Forex_trading_platform_using_Kafka_Storm_Cassandra.html

【讨论】:

  • 您的架构要求没有很好地反映在问题中。
  • 好吧,也许你是对的。以下是我的映射方式:我说,“首先,我们的要求相当简单——写入 HDFS”。这在本节的博客文章中得到了解答:“批量聚合服务的数据管道:慢”。它指出他们使用了“加缪”——这也是我的想法。将来,Kafka/Storm 或 Spark Streaming 会派上用场。无论如何,我的坏!感谢您的宝贵时间。
【解决方案2】:

您可以通过对 DStream 使用 hadoop 操作来保存 DStream 中的数据:

val streamingContext = new StreamingContext(sparkContext, Duration(window))
val tweetStream = TwitterUtils.createStream(streamingContext,...).map(tweet=>tweet.toJSONString)
tweetStream.saveAsTextFiles(pathPrefix, suffix)

假设输入恒定,时间窗口将使您能够控制每个流式传输间隔要处理的消息量。

【讨论】:

  • 我在 JavaStreamingContext 或 StreamingContext 上都没有看到“createTwitterStream”方法。可能它仅在 Scala 中可用?我使用的是 1.1.0 版本的 Spark Streaming。
  • 它叫TwitterUtils.createStream(ssc, ...)我会用确切的电话更新答案。
猜你喜欢
  • 1970-01-01
  • 2018-04-07
  • 2021-03-07
  • 1970-01-01
  • 2015-01-13
  • 2020-03-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多