Spark SQL 中的数据跳过
数据跳过功能可通过根据每个对象关联的摘要元数据,跳过无关的数据对象或文件,从而显著提升 SQL 查询的性能。
数据跳过功能利用开源的 Xskipper 库,通过 Apache Spark 来创建、管理和部署数据跳过索引。 参见 Xskipper——一个可扩展的数据跳过框架。
有关如何使用 Xskipper 的更多详细信息,请参阅:
除了 Xskipper 中的开源功能外,还提供以下功能:
地理空间数据缺失
在使用时空库中的地理空间函数查询地理空间数据集时,您也可以使用数据跳过功能。
- 若要在包含经纬度列的数据集中利用数据跳过功能,您可以为经纬度列建立最小/最大索引。
- 在包含几何列(即 UDT 列)的数据集中,可通过使用内置的 Xskipper 插件来实现数据跳过功能。
以下各节将向您介绍如何使用地理空间插件。
配置地理空间插件
要使用该插件,请使用 "注册" 模块装入相关实现。 请注意,您只能在 IBM Analytics Engine powered by Apache Spark 中的应用程序中使用 Scala ,而不能在 Watson Studio 中使用。
关于 Scala :
import com.ibm.xskipper.stmetaindex.filter.STMetaDataFilterFactory import com.ibm.xskipper.stmetaindex.index.STIndexFactory import com.ibm.xskipper.stmetaindex.translation.parquet.{STParquetMetaDataTranslator, STParquetMetadatastoreClauseTranslator} import io.xskipper._ Registration.addIndexFactory(STIndexFactory) Registration.addMetadataFilterFactory(STMetaDataFilterFactory) Registration.addClauseTranslator(STParquetMetadatastoreClauseTranslator) Registration.addMetaDataTranslator(STParquetMetaDataTranslator)关于 Python :
from xskipper import Xskipper from xskipper import Registration Registration.addMetadataFilterFactory(spark, 'com.ibm.xskipper.stmetaindex.filter.STMetaDataFilterFactory') Registration.addIndexFactory(spark, 'com.ibm.xskipper.stmetaindex.index.STIndexFactory') Registration.addMetaDataTranslator(spark, 'com.ibm.xskipper.stmetaindex.translation.parquet.STParquetMetaDataTranslator') Registration.addClauseTranslator(spark, 'com.ibm.xskipper.stmetaindex.translation.parquet.STParquetMetadatastoreClauseTranslator')
索引建立
要构建索引,可以使用 addCustomIndex API。 请注意,您只能在 IBM Analytics Engine powered by Apache Spark 中的应用程序中使用 Scala ,而不能在 Watson Studio 中使用。
关于 Scala :
import com.ibm.xskipper.stmetaindex.implicits._ // index the dataset val xskipper = new Xskipper(spark, dataset_path) xskipper .indexBuilder() // using the implicit method defined in the plugin implicits .addSTBoundingBoxLocationIndex("location") // equivalent //.addCustomIndex(STBoundingBoxLocationIndex("location")) .build(reader).show(false)关于 Python :
xskipper = Xskipper(spark, dataset_path) # adding the index using the custom index API xskipper.indexBuilder() \ .addCustomIndex("com.ibm.xskipper.stmetaindex.index.STBoundingBoxLocationIndex", ['location'], dict()) \ .build(reader) \ .show(10, False)
支持的功能
支持的地理空间函数列表包括以下内容:
- ST_Distance
- ST_Intersects
- ST_Contains
- ST_Equals
- ST_Crosses
- ST_Touches
- ST_Within
- ST_Overlaps
- ST_EnvelopesIntersect
- ST_IntersectsInterior
加密索引
如果您使用 Parquet 元数据存储,可以选择使用 Parquet 模块化加密 (PME) 对元数据进行加密。 这是通过将元数据本身存储为 Parquet 数据集来实现的,因此可以使用 PME 对它进行加密。 此功能适用于所有输入格式,例如,存储为 CSV 格式的数据集可以使用PME对其元数据进行加密。
在下一节中,除非另有说明,提及页脚、列等内容时,均指元数据对象,而非索引数据集中的对象。
索引加密具有以下模块化和精细化的特点:
- 每个索引都可以加密 (使用每个索引的密钥粒度) 或保留为纯文本
- 页脚 + 对象名称列:
- 元数据对象(其本身是一个 Parquet 文件)的页脚列包含以下内容(除其他内容外):
- 元数据对象的结构图,其中显示了所有已收集索引的类型、参数和列名。 例如,你可以了解到,a 是在列
BloomFilter上定义city的,其假阳性概率为0.1。 - 原始数据集的完整路径,如果是 Hive 元数据存储表,则为表名。
- 元数据对象的结构图,其中显示了所有已收集索引的类型、参数和列名。 例如,你可以了解到,a 是在列
- “对象名称”列存储了所有已建立索引的对象的名称。
- 元数据对象(其本身是一个 Parquet 文件)的页脚列包含以下内容(除其他内容外):
- 页脚 + 元数据列可以是:
两者均使用相同的密钥进行加密。 这是默认设置。 在这种情况下,构成元数据的 Parquet 对象在加密页脚模式下采用明文页脚配置,且对象名称列使用所选密钥进行加密。
两者均为纯文本。 在这种情况下,包含元数据的 Parquet 对象处于纯文本页脚模式,且对象名称列未加密。
如果至少有一个索引被标记为加密状态,则无论是否启用了明文页脚模式,都必须配置页脚密钥。 如果设置了纯文本页脚,则页脚密钥仅用于防篡改。 请注意,在这种情况下,对象名称列不具备防篡改功能。
如果配置了页脚密钥,则至少必须对一个索引进行加密。
在使用索引加密之前,您应查阅 PME 的相关文档,并确保您已熟悉相关概念。
key 时,它始终是标签 "从不密钥本身"。要使用索引加密:
请按照所有步骤操作,以确保已启用 PME。 请参阅 PME (PME)。
执行所有常规的PME配置,包括密钥管理配置。
为数据集创建加密元数据:
- 请按照常规流程创建元数据。
- 配置页脚键。 如果您希望设置一个纯文本页脚 + 对象名称列,请将 设置
io.xskipper.parquet.encryption.plaintext.footer为true(参见下方的示例)。 - 在
IndexBuilder中,对于每个需要加密的索引,请添加该索引所用密钥的标签。
若要在查询时使用元数据或刷新现有元数据,除确保密钥可访问所需的常规 PME 配置外(实际上与读取加密数据集所需的配置完全相同),无需进行其他设置。
样本
以下示例演示了如何使用名为 k1 的键作为页脚 + 对象名称键,使用名为 k2 的键对 进行加密 MinMax , temp同时为 创建一个 ValueList ,该 city 保留为明文。 请注意,您只能在 IBM Analytics Engine powered by Apache Spark 中的应用程序中使用 Scala ,而不能在 Watson Studio 中使用。
关于 Scala :
// index the dataset val xskipper = new Xskipper(spark, dataset_path) // Configuring the JVM wide parameters val jvmComf = Map( "io.xskipper.parquet.mdlocation" -> md_base_location, "io.xskipper.parquet.mdlocation.type" -> "EXPLICIT_BASE_PATH_LOCATION") Xskipper.setConf(jvmConf) // set the footer key val conf = Map( "io.xskipper.parquet.encryption.footer.key" -> "k1") xskipper.setConf(conf) xskipper .indexBuilder() // Add an encrypted MinMax index for temp .addMinMaxIndex("temp", "k2") // Add a plaintext ValueList index for city .addValueListIndex("city") .build(reader).show(false)适用于 Python
xskipper = Xskipper(spark, dataset_path) # Add JVM Wide configuration jvmConf = dict([ ("io.xskipper.parquet.mdlocation", md_base_location), ("io.xskipper.parquet.mdlocation.type", "EXPLICIT_BASE_PATH_LOCATION")]) Xskipper.setConf(spark, jvmConf) # configure footer key conf = dict([("io.xskipper.parquet.encryption.footer.key", "k1")]) xskipper.setConf(conf) # adding the indexes xskipper.indexBuilder() \ .addMinMaxIndex("temp", "k1") \ .addValueListIndex("city") \ .build(reader) \ .show(10, False)
如果你希望页脚和对象名称保持为纯文本模式(如上所述),则需要添加以下配置参数:
关于 Scala :
// index the dataset val xskipper = new Xskipper(spark, dataset_path) // Configuring the JVM wide parameters val jvmComf = Map( "io.xskipper.parquet.mdlocation" -> md_base_location, "io.xskipper.parquet.mdlocation.type" -> "EXPLICIT_BASE_PATH_LOCATION") Xskipper.setConf(jvmConf) // set the footer key val conf = Map( "io.xskipper.parquet.encryption.footer.key" -> "k1", "io.xskipper.parquet.encryption.plaintext.footer" -> "true") xskipper.setConf(conf) xskipper .indexBuilder() // Add an encrypted MinMax index for temp .addMinMaxIndex("temp", "k2") // Add a plaintext ValueList index for city .addValueListIndex("city") .build(reader).show(false)适用于 Python
xskipper = Xskipper(spark, dataset_path) # Add JVM Wide configuration jvmConf = dict([ ("io.xskipper.parquet.mdlocation", md_base_location), ("io.xskipper.parquet.mdlocation.type", "EXPLICIT_BASE_PATH_LOCATION")]) Xskipper.setConf(spark, jvmConf) # configure footer key conf = dict([("io.xskipper.parquet.encryption.footer.key", "k1"), ("io.xskipper.parquet.encryption.plaintext.footer", "true")]) xskipper.setConf(conf) # adding the indexes xskipper.indexBuilder() \ .addMinMaxIndex("temp", "k1") \ .addValueListIndex("city") \ .build(reader) \ .show(10, False)
连接操作中的数据跳过(仅限 Spark 3)
在 Spark 3 中,您可以在连接查询中使用数据跳过功能,例如:
SELECT *
FROM orders, lineitem
WHERE l_orderkey = o_orderkey and o_custkey = 800
此示例展示了一个基于 TPC-H 基准模式(参见 TPC-H )的星型模式,其中 lineitem 是事实表,包含大量记录;而 orders 表是维度表,其记录数量与事实表相比相对较少。
上述查询对 orders 表设置了谓词,而该表包含的记录数量较少,这意味着使用 min/max 操作时,数据跳过带来的效益不大。
动态数据跳过是一项功能,它使上述查询能够利用数据跳过机制:首先根据表 orders 上的条件提取相关 l_orderkey 值,然后利用这些值将谓词下推至 l_orderkey 该表,并利用支持数据跳过的索引过滤掉无关对象。
要使用此功能,请启用以下优化规则。 请注意,您只能在 IBM Analytics Engine powered by Apache Spark 中的应用程序中使用 Scala ,而不能在 Watson Studio 中使用。
关于 Scala :
import com.ibm.spark.implicits. spark.enableDynamicDataSkipping()关于 Python :
from sparkextensions import SparkExtensions SparkExtensions.enableDynamicDataSkipping(spark)
然后像往常一样使用 Xskipper API,您的查询将受益于数据跳过功能。
例如,在上述查询中,使用 min/max l_orderkey 建立索引可以跳过该表 lineitem ,从而提升查询性能。
对旧版元数据的支持
Xskipper 能够无缝支持由 MetaIndexManager 生成的旧版元数据。 旧的元数据可用于跳过,因为对 Xskipper 元数据的更新将在下一次刷新操作中自动完成。
在列出索引或执行 操作时,如果看到 DEPRECATED_SUPPORTED 索引前带有 标记 describeIndex ,则表示该元数据版本已弃用,但仍受支持,跳过该索引也是可行的。 下次刷新操作将自动更新元数据。