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 元数据存储表,则为表名。
    • “对象名称”列存储了所有已建立索引的对象的名称。
  • 页脚 + 元数据列可以是:
    • 两者均使用相同的密钥进行加密。 这是默认设置。 在这种情况下,构成元数据的 Parquet 对象在加密页脚模式下采用明文页脚配置,且对象名称列使用所选密钥进行加密。

    • 两者均为纯文本。 在这种情况下,包含元数据的 Parquet 对象处于纯文本页脚模式,且对象名称列未加密。

      如果至少有一个索引被标记为加密状态,则无论是否启用了明文页脚模式,都必须配置页脚密钥。 如果设置了纯文本页脚,则页脚密钥仅用于防篡改。 请注意,在这种情况下,对象名称列不具备防篡改功能。

      如果配置了页脚密钥,则至少必须对一个索引进行加密。

在使用索引加密之前,您应查阅 PME 的相关文档,并确保您已熟悉相关概念。

重要信息: 使用索引加密时,每当在任何 Xskipper API 中配置 key 时,它始终是标签 "从不密钥本身"。

要使用索引加密:

  1. 请按照所有步骤操作,以确保已启用 PME。 请参阅 PME (PME)

  2. 执行所有常规的PME配置,包括密钥管理配置。

  3. 为数据集创建加密元数据:

    1. 请按照常规流程创建元数据。
    2. 配置页脚键。 如果您希望设置一个纯文本页脚 + 对象名称列,请将 设置 io.xskipper.parquet.encryption.plaintext.footertrue (参见下方的示例)。
    3. IndexBuilder中,对于每个需要加密的索引,请添加该索引所用密钥的标签。

    若要在查询时使用元数据或刷新现有元数据,除确保密钥可访问所需的常规 PME 配置外(实际上与读取加密数据集所需的配置完全相同),无需进行其他设置。

样本

以下示例演示了如何使用名为 k1 的键作为页脚 + 对象名称键,使用名为 k2 的键对 进行加密 MinMaxtemp同时为 创建一个 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 ,则表示该元数据版本已弃用,但仍受支持,跳过该索引也是可行的。 下次刷新操作将自动更新元数据。