Spark-Phoenix如何按条件拉取HBase指定时段数据避免全量拉取
答复
针对你使用phoenix-4.14.1-HBase-1.3版本的场景,两个能力均原生支持,不需要全量拉取HBase表数据,具体说明如下:
- 关于指定SQL拉取指定范围数据:和Spark对接普通JDBC数据源的逻辑完全一致,spark-phoenix connector支持直接传入自定义查询SQL,在Phoenix侧完成数据过滤后再返回结果给Spark。
读取时直接通过phoenix.sql选项传入你写好的带1小时时间范围过滤的SQL即可,示例代码如下:
如果你的时间字段提前建好了Phoenix二级索引,这个查询的拉取速度会有明显提升。val hourRangeDf = spark.read .format("org.apache.phoenix.spark") .option("zkUrl", "集群ZooKeeper连接地址") .option("phoenix.sql", "select 所需字段 from 你的业务表 where 时间字段 >= 起始时间值 and 时间字段 < 结束时间值") .load() - 关于谓词下推能力:spark-phoenix connector原生实现了Spark DataSource的谓词下推、列裁剪规则,你直接加载Phoenix表注册临时视图后写SparkSQL加where条件过滤,connector会自动把可适配的过滤条件下推到Phoenix侧执行,不会拉取全表到Spark内存后再做过滤。
你可以通过执行df.explain(true)查看物理执行计划确认下推状态,如果计划中PushedFilters列表里的时间过滤条件前带*标记,就代表对应条件已经成功下推。需要注意如果你的时间字段是自定义编码的字符串/数值类型,要保证Spark侧的字段类型和Phoenix侧字段类型匹配,否则可能出现下推失效的问题。
内容的提问来源于stack exchange,提问作者cnidaye
相关产品推荐
相关产品推荐

