Data skipping para Spark SQL

O data skipping pode impulsionar significativamente o desempenho de consultas SQL ignorando arquivos ou objetos de dados irrelevantes com base em um resumo dos metadados associados a cada objeto.

O data skipping usa a biblioteca Xskipper de software livre para criar, gerenciar e implementar índices de data skipping com o Apache Spark. Consulte Xskipper - Uma estrutura extensível de data skipping.

Para obter mais detalhes sobre como trabalhar com o Xskipper, consulte:

Além dos recursos de software livre no Xskipper, os recursos a seguir também estão disponíveis:

Data skipping geoespacial

Também pode-se ignorar dados ao consultar conjuntos de dados geoespaciais usando funções geoespaciais a partir da biblioteca relativa ao espaço e tempo.

  • Para se beneficiar do data skipping em conjuntos de dados com colunas de latitude e longitude, é possível coletar os índices mín-máx nas colunas de latitude e longitude.
  • É possível ignorar dados em conjuntos de dados com uma coluna de geometria (uma coluna UDT) usando um Plug-in Xskipper integrado.

As seções a seguir mostram como trabalhar com o plug-in geoespacial.

Configurando o plug-in geoespacial

Para usar o plugin, carregue as implementações relevantes usando o módulo Registro. Observe que você só pode usar Scala em aplicativos em IBM Analytics Engine powered by Apache Spark, não em watsonx.ai Studio.

  • Para o 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)
    
  • Para 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')
    

Construção de índice

Para construir um índice, é possível usar a API do addCustomIndex Observe que você só pode usar Scala em aplicativos em IBM Analytics Engine powered by Apache Spark, não em watsonx.ai Studio.

  • Para o 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)
    
  • Para 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)
    

Funções suportadas do

A lista de funções geoespaciais suportadas inclui o seguinte:

  • ST_Distance
  • ST_Intersects
  • ST_Contains
  • ST_Equals
  • ST_Crosses
  • ST_Touches
  • ST_Within
  • ST_Overlaps
  • ST_EnvelopesIntersect
  • ST_IntersectsInterior

Criptografia de índices

Se um armazenamento de metadados Parquet for usado, opcionalmente, os metadados poderão ser criptografados usando Parquet Modular Encryption (PME). Isso é feito armazenando os próprios metadados como um conjunto de dados do Parquet e, sendo assim, o PME pode ser usado para criptografá-los. Esse recurso se aplica a todos os formatos de entrada, por exemplo, um conjunto de dados armazenado em formato CSV pode ter seus metadados criptografados usando o PME.

Na seção a seguir, salvo se especificado de outra forma, as referências a rodapés, colunas etc., serão a respeito dos objetos de metadados, e não dos objetos no conjunto de dados indexado.

A criptografia de índice é modular e granular conforme mostrado a seguir:

  • Cada índice pode ser criptografado (com uma granularidade de chave por índice) ou deixado em texto simples
  • Coluna de rodapé + nome de objeto:
    • A coluna de rodapé do objeto de metadados, que é um arquivo do Parquet em si, contém, entre outras coisas:
      • Esquema do objeto de metadados, que revela os tipos, parâmetros e nomes de coluna para todos os índices coletados. Por exemplo, você pode aprender que um BloomFilter é definido na coluna city com uma probabilidade de falso-positivo igual a 0.1.
      • Caminho completo para o conjunto de dados original ou o nome da tabela em caso de uma tabela do metastore Hive.
    • A coluna de nome de objeto armazena os nomes de todos os objetos indexados.
  • A coluna de rodapé + metadados pode ser:
    • Criptografada usando a mesma chave. Este é o padrão. Nesse caso, a configuração do rodapé de texto sem formatação para os objetos do Parquet contendo os metadados é para o modo de rodapé criptografado e a coluna de nome de objeto é criptografada usando a chave selecionada.

    • Texto sem formatação. Nesse caso, os objetos do Parquet contendo os metadados estão no modo de rodapé de texto sem formatação e a coluna de nome de objeto não é criptografada.

      Se pelo menos um índice for marcado como criptografado, uma chave de rodapé deverá ser configurada, independentemente de o modo de rodapé de texto sem formatação estar ou não ativado. Se o rodapé de texto sem formatação estiver configurado, a chave de rodapé será usada para proteção contra violação. Observe que, nesse caso, a coluna de nome de objeto não tem proteção contra violação.

      Se uma chave de rodapé estiver configurada, pelo menos um índice deverá ser criptografado.

Antes de usar a criptografia de índice, deve-se verificar a documentação no PME e certificar-se de que você esteja familiarizado com os conceitos.

Importante: ao usar a criptografia de índice, sempre que um key é configurado em qualquer API do Xskipper, ele sempre é o rótulo ` NUNCA a chave em si `.

Para usar criptografia de índice:

  1. Siga todas as etapas para assegurar que o PME esteja ativado. Consulte PME.

  2. Execute todas as configurações do PME regular, incluindo configurações de Gerenciamento de chaves.

  3. Crie metadados criptografados para um conjunto de dados:

    1. Siga o fluxo regular para criar metadados.
    2. Configure uma chave de rodapé. Se quiser configurar um rodapé de texto simples + coluna nome de objeto, configure io.xskipper.parquet.encryption.plaintext.footer para true (Veja amostras abaixo).
    3. Em IndexBuilder, para cada índice que deseja criptografar, inclua o rótulo da chave para usar para esse índice.

    Para usar metadados durante o tempo de consulta ou para atualizar os metadados existentes, nenhuma configuração é necessária além da configuração de PME regular, exigida para garantir que as chaves estejam acessíveis (literalmente a mesma configuração necessária para ler um conjunto de dados criptografados).

Amostras

As amostras a seguir exibem a criação de metadados usando uma chave denominada k1 como uma chave de rodapé + nome de objeto e uma chave denominada k2 como uma chave para criptografar um MinMax para temp, enquanto também cria um ValueList para city, que é deixado em texto simples. Observe que você só pode usar Scala em aplicativos em IBM Analytics Engine powered by Apache Spark, não em watsonx.ai Studio.

  • Para o 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)
    
  • Para 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)
    

Se quiser que o rodapé + nome do objeto sejam deixados no modo de texto sem formatação (conforme mencionado acima), será necessário incluir o parâmetro de configuração:

  • Para o 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)
    
  • Para 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)
    

Data skipping com junções (apenas para Spark 3)

Com o Spark 3, é possível usar o data skipping em consultas de junção como:

SELECT *
FROM orders, lineitem 
WHERE l_orderkey = o_orderkey and o_custkey = 800

Este exemplo mostra um esquema de estrela com base no esquema de benchmark TPC-H (consulte TPC-H) em que item de linha é uma tabela de fatos e contém muitos registros, enquanto a tabela de pedidos é uma tabela de dimensões que possui um número relativamente pequeno de registros em comparação com as tabelas de fatos.

A consulta acima conta com um predicado na tabela de pedidos que contém um pequeno número de registros, o que significa que o uso de mín-máx não será tão beneficiado pelo data skipping.

Ignorar dados dinamicamente é um recurso que possibilita consultas como a acima, para se beneficiar de dados ignorando primeiro por extração dos valores l_orderkey relevantes com base na condição na tabela orders e, em seguida, utilizá-los para empurrar para baixo um predicado no l_orderkey que usa índices de ignorando dados para filtrar objetos irrelevantes.

Para usar esse recurso, ative a regra de otimização a seguir Observe que você só pode usar Scala em aplicativos em IBM Analytics Engine powered by Apache Spark, não em watsonx.ai Studio.

  • Para o Scala:

      import com.ibm.spark.implicits.
    
      spark.enableDynamicDataSkipping()
    
  • Para Python:

        from sparkextensions import SparkExtensions
    
        SparkExtensions.enableDynamicDataSkipping(spark)
    

Em seguida, use a API da Xskipper normalmente e suas consultas serão beneficiadas pelo uso do data skipping.

Por exemplo, na consulta acima, a indexação l_orderkey usando mín. / máx. irá ativar a passagem sobre a tabela lineitem e melhorará o desempenho da consulta.

Suporte para metadados mais antigos

A Xskipper suporta os metadados mais antigos criados facilmente pelo MetaIndexManager. Os metadados mais antigos podem ser usados para skipping, já que as atualizações nos metadados da Xskipper são realizadas automaticamente pela operação de atualização seguinte.

Se você vir a DEPRECATED_SUPPORTED na frente de um índice ao listar índices ou executar uma operação describeIndex, a versão de metadados será descontinuada, mas ainda será suportada e a opção de ignorar funcionará. A próxima operação de atualização irá atualizar os metadados automaticamente.