【问题标题】:How can I read one tables partial data from hbase with pyspark and shc(spark hbase connector) instead of whole dataset?如何使用 pyspark 和 shc(spark hbase 连接器)而不是整个数据集从 hbase 读取一个表的部分数据?
【发布时间】:2019-07-20 05:03:13
【问题描述】:

我正在使用 pyspark 通过 shc 访问 hbase 的表。该表有大量记录,但我的spark集群只有三台服务器,性能很差。我认为从那个hbase表中读取整个数据然后用spark的过滤器处理它是不合理的,那么我怎么能用pyspark和shc从hbase中读取部分数据呢? 例如,我想要筛选具有起始值结束值或筛选列的行键

有一个基本的读写方法,谢谢

from pyspark.sql import SparkSession
spark = SparkSession.builder.master('localhost').appName('test_1').getOrCreate()

def test_shc():
    catalog = ''.join("""{
      "table":{"namespace":"test", "name":"test_shc"},
      "rowkey":"key",
      "columns":{
      "col0":{"cf":"rowkey", "col":"key", "type":"string"},
      "col1":{"cf":"result", "col":"class", "type":"string"}
      }
      }""".split())

    data_source_format = 'org.apache.spark.sql.execution.datasources.hbase'
    df = spark.sparkContext.parallelize([('a', '1.0'), ('b', '2.0')]).toDF(schema=['col0', 'col1'])
    df.show()
    df.write.options(catalog=catalog, newTable="5").format(data_source_format).save()
    df_read = spark.read.options(catalog=catalog).format(data_source_format).load()
    df_read.show()

【问题讨论】:

    标签: python pyspark hbase


    【解决方案1】:

    使用

    spark.read.options(catalog=catalog).format(data_source_format).load().limit(n)

    在加载数据时。 limit(n) 将限制读取的记录数量。

    【讨论】:

    • 谢谢,但是如何实现获取一系列rowkey数据的功能呢?如果 rowkey 像 '20190804xxxxx',我想在给定的时间范围内获取数据,例如 (20190705, 20190803]
    【解决方案2】:

    从那时起,您可能已经找到了答案...但是,您可以避免加载整个文件的方式,但只有您想要的行就是您所说的,换句话说,根据行键应用过滤器之后正在加载。
    这样 SHC 不会加载整个表,因为过滤器将在 HBase 上下推,只有过滤器中的数据会返回到 Spark,如您在此处看到的:https://github.com/hortonworks-spark/shc/issues/108,其中

    def withCatalog(cat: String): DataFrame = {
      sqlContext 
      .read
      .options(Map(HBaseTableCatalog.tableCatalog->cat))
      .format("org.apache.spark.sql.execution.datasources.hbase")
      .load()
     }
    

    请注意,过滤非行键列会导致对表进行全面扫描。在这种情况下,查询效率不是很高,一种建议的方法是使用 SingleColumnValueFilter 将计算推送到 HBase 层,以便在将数据传输到 Spark 之前尽可能多地进行过滤。

    更多详情:https://www.ifi.uzh.ch/dam/jcr:fbf88795-65a3-4bc1-9a8a-60548d7f2f91/SHC%20Distributed%20Query%20Processing.pdf

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-01-26
      相关资源
      最近更新 更多