【问题标题】:How to use Apache spark as Query Engine?如何使用 Apache spark 作为查询引擎?
【发布时间】:2017-01-22 06:12:04
【问题描述】:

我正在使用 Apache Spark 进行大数据处理。数据从平面文件源或 JDBC 源加载到数据帧。 Job 是使用 spark sql 从数据框中搜索特定记录。

所以我必须一次又一次地运行工作以获取新的搜索词。每次我必须使用 spark submit 提交 Jar 文件以获得结果。 由于数据大小为 40.5 GB,每次将相同的数据重新加载到数据帧以获取不同查询的结果会变得很繁琐

所以我需要的是,

  • 如果我可以一次加载数据帧中的数据并多次查询它而无需多次提交 jar 的方法?
  • 如果我们可以使用 spark 作为搜索引擎/查询引擎?
  • 如果我们可以将数据加载到数据框中一次并使用RestAP远程查询数据框

> My Spark Deployment 的当前配置是

  • 5 节点集群。
  • 在纱线 rm 上运行。

我曾尝试使用 spark-job 服务器,但它每次都会运行该作业。

【问题讨论】:

  • 如果我们可以使用 Spark sql 使用 Rest Api 查询现有数据框? - 是的。 如果我们可以使用 spark 作为搜索引擎/查询引擎? - 基于意见,但有 40GB 的数据,只需使用像样的 RDBMS。您将获得更好的投资回报率。 如果我可以加载数据框一次的方法 - 不止一个。从内置节俭服务器到不同的休息选项和数据网格。
  • @zero323 会很好。能不能解释的更准确一点?
  • @KamalPradhan 应该可以使用 spark-jobserver 在作业之间缓存 RDD。我认为你必须命名你的 RDD 才能工作。详情见here
  • @GrahamS 问题是我们必须使用 spark-sql 查询数据框。因此,如果我们可以一次将数据加载到数据框中并使用 RestAPI 远程查询数据框,我的问题将得到解决。

标签: apache-spark apache-spark-sql spark-streaming spark-dataframe


【解决方案1】:

您可能对HiveThriftServer 和 Spark 集成感兴趣。

基本上你启动一个 Hive Thrift 服务器并从 SparkContext 注入你的 HiveContext 构建:

...
val sql = new HiveContext(sc)
sql.setConf("hive.server2.thrift.port", "10001")
...
dataFrame.registerTempTable("myTable")
HiveThriftServer2.startWithContext(sql)
...

有几个客户端库和工具可以查询服务器: https://cwiki.apache.org/confluence/display/Hive/HiveServer2+Clients

包括 CLI 工具 - beeline

参考: https://medium.com/@anicolaspp/apache-spark-as-a-distributed-sql-engine-4373e254e0f9#.3ntbhdxvr

【讨论】:

  • 感谢您的帮助。我们可以使用 rest api 或一些 web 服务查询服务器吗?它是否也能够在纱线集群模式下进行扩展。
  • 您有专门的 java (JDBC)、python 和 ruby​​ 库可以查询此服务 - 在响应中提到 - “客户端”链接
  • 也请查看官方文档:spark.apache.org/docs/latest/…
【解决方案2】:

您还可以使用 spark+kafka 流式集成。只是您必须通过 kafka 发送查询,以便流式 API 接收。如果它简单的话,那是一种在市场上迅速流行起来的设计模式。

  1. 在您的查找数据上创建数据集。

  2. 通过 Kafka 启动 Spark 流式查询。

  3. 从您的 Kafka 主题中获取 sql

  4. 对已创建的数据集执行查询

这应该照顾你的用例。

希望这会有所帮助!

【讨论】:

    【解决方案3】:

    对于 spark 搜索引擎,如果您需要全文搜索功能和/或文档级别评分 - 并且您没有弹性搜索基础设施 - 您可以尝试 Spark Search - 它为 @带来了 Apache Lucene 支持987654323@.

    df.rdd.searchRDD().save("/tmp/hdfs-pathname")
    val restoredSearchRDD: SearchRDD[Person] = SearchRDD.load[Person](sc, "/tmp/hdfs-pathname")
    restoredSearchRDD.searchList("(fistName:Mikey~0.8) OR (lastName:Wiliam~0.4) OR (lastName:jonh~0.2)",
                                topKByPartition = 10)
                     .map(doc => s"${doc.source.firstName}=${doc.score}"
                     .foreach(println)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2015-10-08
      • 2015-05-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-04-18
      相关资源
      最近更新 更多