【问题标题】:How to view pyspark temporary tables on Thrift server?如何在 Thrift 服务器上查看 pyspark 临时表?
【发布时间】:2019-11-11 05:30:43
【问题描述】:

我正在尝试通过 Thrift 在 pyspark 上创建一个临时表。我的最终目标是能够使用 JDBC 从 DBeaver 等数据库客户端访问它。

我首先使用 beeline 进行测试。

这就是我正在做的。

  1. 在我自己的机器上使用 docker 启动了一个集群,并在 spark-defaults.conf 上添加了 spark.sql.hive.thriftServer.singleSession true
  2. 启动 Pyspark shell(为了测试)并运行以下代码:

    from pyspark.sql import Row l = [('Ankit',25),('Jalfaizy',22),('saurabh',20),('Bala',26)] rdd = sc.parallelize(l) people = rdd.map(lambda x: Row(name=x[0], age=int(x[1]))) people = people.toDF().cache() peebs = people.createOrReplaceTempView('peebs') result = sqlContext.sql('select * from peebs')

    到目前为止一切顺利,一切正常。

  3. 在另一个终端上,我初始化了 spark thrift 服务器: ./sbin/start-thriftserver.sh --hiveconf hive.server2.thrift.port=10001 --conf spark.executor.cores=1 --master spark://172.18.0.2:7077

    服务器似乎正常启动,我可以看到 pyspark 和 thrift 服务器作业在我的 spark 集群主 UI 上运行。

  4. 然后我使用直线连接到集群

    ./bin/beeline beeline> !connect jdbc:hive2://172.18.0.2:10001

    这就是我得到的

    连接到 jdbc:hive2://172.18.0.2:10001
    输入 jdbc 的用户名:hive2://172.18.0.2:10001:
    输入 jdbc:hive2://172.18.0.2:10001 的密码:
    2019-06-29 20:14:25 INFO Utils:310 - 提供的权限:172.18.0.2:10001
    2019-06-29 20:14:25 INFO Utils:397 - 已解决的权限:172.18.0.2:10001
    2019-06-29 20:14:25 信息 HiveConnection:203 - 将尝试使用 JDBC Uri 打开客户端传输:jdbc:hive2://172.18.0.2:10001
    连接到:Spark SQL(2.3.3 版)
    驱动程序:Hive JDBC(版本 1.2.1.spark2)
    事务隔离:TRANSACTION_REPEATABLE_READ

    好像没问题。

  5. 当我列出 show tables; 时,我什么都看不到。

我想强调两件有趣的事:

  1. 当我启动 pyspark 时,我会收到这些警告

    WARN ObjectStore:6666 - 在 Metastore 中找不到版本信息。 hive.metastore.schema.verification 未启用,因此记录架构版本 1.2.0

    WARN ObjectStore:568 - 获取数据库默认值失败,返回 NoSuchObjectException

    WARN ObjectStore:568 - 获取数据库 global_temp 失败,返回 NoSuchObjectException

  2. 当我启动 thrift 服务器时,我得到了这些:

    rsync 来自 spark://172.18.0.2:7077
    ssh:无法解析主机名 spark:名称或服务未知
    rsync:连接意外关闭(到目前为止已收到 0 个字节)[接收方]
    rsync 错误:在 io.c(235) [Receiver=3.1.2] 出现无法解释的错误(代码 255)
    正在启动 org.apache.spark.sql.hive.thriftserver.HiveThriftServer2,登录到 ...

我经历了好几篇帖子和讨论。我看到有人说我们不能通过 thrift 公开临时表,除非您从同一代码中启动服务器。如果这是真的,我该如何在 python (pyspark) 中做到这一点?

谢谢

【问题讨论】:

    标签: python apache-spark pyspark thrift spark-thriftserver


    【解决方案1】:

    如果有人需要在 Spark Streaming 中执行此操作,我可以让它像这样工作。

    from pyspark.sql.functions import *
    from pyspark.sql.types import *
    from pyspark import SparkContext, SparkConf
    from pyspark.sql import SparkSession
    from py4j.java_gateway import java_import
    
    spark = SparkSession \
        .builder \
        .appName('restlogs_qlik') \
        .enableHiveSupport()\
        .config('spark.sql.hive.thriftServer.singleSession', True)\
        .config('hive.server2.thrift.port', '10001') \
        .getOrCreate()
    
    sc=spark.sparkContext
    sc.setLogLevel('INFO')
    
    #Order matters! 
    java_import(sc._gateway.jvm, "")
    sc._gateway.jvm.org.apache.spark.sql.hive.thriftserver.HiveThriftServer2.startWithContext(spark._jwrapped)
    
    #Define schema of json
    schema = StructType().add("partnerid", "string").add("sessionid", "string").add("functionname", "string").add("functionreturnstatus", "string")
    
    #load data into spark-structured streaming
    df = spark \
          .readStream \
          .format("kafka") \
          .option("kafka.bootstrap.servers", "localhost:9092") \
          .option("subscribe", "rest_logs") \
          .load() \
          .select(from_json(col("value").cast("string"), schema).alias("parsed_value"))
    
    #Print output
    query = df.writeStream \
                .outputMode("append") \
                .format("memory") \
                .queryName("view_figures") \
                .start()
    
    query.awaitTermination();
    

    在你启动它之后,你可以用蜂巢的 JDBC 来测试。我无法理解的是我必须在同一个脚本中启动 Thrift 服务器。这是启动脚本的方法。

        spark-submit --master local[2] \
    --conf "spark.driver.extraClassPath=D:\\Libraries\\m2_repository\\org\\apache\\kafka\\kafka-clients\\2.0.0\\kafka-clients-2.0.0.jar" \
    --packages org.apache.spark:spark-sql-kafka-0-10_2.11:2.4.1 \
    "C:\\Temp\\spark_kafka.py"
    

    希望这对某人有所帮助。顺便说一句,我在初步研究阶段,所以不要评判我。

    【讨论】:

      【解决方案2】:

      createOrReplaceTempView 创建一个内存表。 Spark thrift 服务器需要在我们创建内存表的同一驱动程序 JVM 上启动。
      在上面的例子中,创建表的驱动和运行STS(Spark Thrift server)的驱动是不同的。
      两种选择
      1.在启动STS的同一JVM中使用createOrReplaceTempView创建表。
      2.使用后备元存储,并使用org.apache.spark.sql.DataFrameWriter#saveAsTable创建表,以便独立于JVM访问表(实际上没有任何Spark驱动程序。

      关于错误:
      1. 与客户端和服务器元存储版本有关。
      2. 好像一些 rsync 脚本试图解码 spark:\\ url
      两者似乎都与问题无关。

      【讨论】:

      • 关于#1,有没有一种方法可以在同一个 JVM 中创建视图,而无需在代码末尾启动 STS?我的意思是,有没有办法创建它并仍然使用 start-thriftserver.sh 启动 STS?你有一些示例代码吗?关于 #2 的问题是,我想在通过 JDBC 公开之前对数据进行一些转换,但我不想再生成或处理一个表。
      • @BrunoFaria,你可以使用org.apache.spark.sql.hive.thriftserver.HiveThriftServer2#start。在驱动程序代码中,创建临时表,并在驱动程序代码的最后一行启动 STS。还需要正确设置节俭属性(检查cwiki.apache.org/confluence/display/Hive/… 以获取前缀为hive.server2 的属性)
      【解决方案3】:

      在进行了几次测试后,我想出了一个适合我的简单(无身份验证)代码。

      请务必注意,如果您想通过 JDBC 提供临时表,您需要在同一个 JVM(同一个 spark 作业)中启动 thrift 服务器,并确保代码挂起,以便应用程序继续在集群中运行。

      按照我创建的工作示例代码供参考:

      import time
      from pyspark import SparkContext, SparkConf
      from pyspark.sql import SparkSession
      from py4j.java_gateway import java_import
      
      spark = SparkSession \
          .builder \
          .appName('the_test') \
          .enableHiveSupport()\
          .config('spark.sql.hive.thriftServer.singleSession', True)\
          .config('hive.server2.thrift.port', '10001') \
          .getOrCreate()
      
      sc=spark.sparkContext
      sc.setLogLevel('INFO')
      
      java_import(sc._gateway.jvm, "")
      
      
      from pyspark.sql import Row
      l = [('John', 20), ('Heather', 34), ('Sam', 23), ('Danny', 36)]
      rdd = sc.parallelize(l)
      people = rdd.map(lambda x: Row(name=x[0], age=int(x[1])))
      people = people.toDF().cache()
      peebs = people.createOrReplaceTempView('peebs')
      
      sc._gateway.jvm.org.apache.spark.sql.hive.thriftserver.HiveThriftServer2.startWithContext(spark._jwrapped)
      
      while True:
          time.sleep(10)
      

      我只是在我的 spark-submit 中使用了上面的 .py,并且我能够通过 beeline 通过 JDBC 进行连接,并使用 Hive JDBC 驱动程序使用 DBeaver。

      【讨论】:

        猜你喜欢
        • 2022-09-23
        • 2012-11-21
        • 2017-10-03
        • 1970-01-01
        • 2017-03-14
        • 1970-01-01
        • 1970-01-01
        • 2018-06-05
        • 1970-01-01
        相关资源
        最近更新 更多