Hinweis
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, sich anzumelden oder das Verzeichnis zu wechseln.
Für den Zugriff auf diese Seite ist eine Autorisierung erforderlich. Sie können versuchen, das Verzeichnis zu wechseln.
Hyperspace bietet Apache Spark-Benutzern die Möglichkeit, Indizes für ihre Datasets (z. B. CSV, JSON und Parquet) zu erstellen und diese zur potenziellen Beschleunigung von Abfragen und Workloads zu nutzen.
In diesem Artikel werden die Grundlagen von Hyperspace vorgestellt, wobei die Einfachheit hervorgehoben und gezeigt wird, wie es von nahezu allen Personen verwendet werden kann.
Haftungsausschluss: Hyperspace unterstützt Sie unter zwei Umständen beim Beschleunigen Ihrer Workloads oder Abfragen:
- Die Abfragen enthalten Filter für Prädikate mit hoher Selektivität. Beispielsweise können Sie 100 übereinstimmende Zeilen aus einer Million Kandidatenzeilen auswählen.
- Die Abfragen enthalten einen Join, der umfassende Shufflevorgänge erfordert. Beispielsweise können Sie ein 100-GB-Dataset mit einem 10-GB-Dataset verknüpfen.
Sie sollten Ihre Workloads sorgfältig überwachen und von Fall zu Fall bestimmen, ob die Indizierung für Sie hilfreich ist.
Dieses Dokument liegt für Python, C# und Scala auch in Notebookform vor.
Einrichten
Hinweis
Hyperspace wird in Azure Synapse Runtime for Apache Spark 3.1 (nicht unterstützt) und Azure Synapse Runtime for Apache Spark 3.2 (Supportende angekündigt) unterstützt. Beachten Sie jedoch, dass Hyperspace in Azure Synapse Runtime für Apache Spark 3.3 (GA) nicht unterstützt wird.
Beginnen Sie mit einer neuen Spark-Sitzung. Da es sich bei diesem Dokument um ein Tutorial handelt, das lediglich veranschaulichen soll, was Hyperspace bieten kann, nehmen Sie eine Konfigurationsänderung vor, durch die hervorgehoben werden kann, wie Hyperspace bei kleinen Datasets vorgeht.
Standardmäßig verwendet Spark Broadcastjoins zur Optimierung von Joinabfragen, wenn die Datengröße für eine Joinseite klein ist (was bei den Beispieldaten, die wir in diesem Tutorial verwenden, der Fall ist). Daher deaktivieren wir Broadcastjoins, sodass Spark später beim Ausführen von Joinabfragen entsprechende Sort-Merge-Joins verwendet. Dies dient vor allem dazu zu zeigen, wie Hyperspace-Indizes im großen Stil zur Beschleunigung von Joinabfragen verwendet würden.
Die Ausgabe der Ausführung der nachfolgenden Zelle zeigt einen Verweis auf die erfolgreich erstellte Spark-Sitzung und gibt „–1“ als Wert für die geänderte Joinkonfiguration aus. Damit wird angezeigt, dass der Broadcastjoin erfolgreich deaktiviert wurde.
// Start your Spark session
spark
// Disable BroadcastHashJoin, so Spark will use standard SortMergeJoin. Currently, Hyperspace indexes utilize SortMergeJoin to speed up query.
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1)
// Verify that BroadcastHashJoin is set correctly
println(spark.conf.get("spark.sql.autoBroadcastJoinThreshold"))
# Start your Spark session.
spark
# Disable BroadcastHashJoin, so Spark will use standard SortMergeJoin. Currently, Hyperspace indexes utilize SortMergeJoin to speed up query.
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1)
# Verify that BroadcastHashJoin is set correctly
print(spark.conf.get("spark.sql.autoBroadcastJoinThreshold"))
// Disable BroadcastHashJoin, so Spark will use standard SortMergeJoin. Currently, Hyperspace indexes utilize SortMergeJoin to speed up query.
spark.Conf().Set("spark.sql.autoBroadcastJoinThreshold", -1);
// Verify that BroadcastHashJoin is set correctly.
Console.WriteLine(spark.Conf().Get("spark.sql.autoBroadcastJoinThreshold"));
Ergebnis:
res3: org.apache.spark.sql.SparkSession = org.apache.spark.sql.SparkSession@297e957d
-1
Datenvorbereitung
Zur Vorbereitung Ihrer Umgebung erstellen Sie Beispieldatensätze, die Sie als Parquet-Datendateien speichern. Parquet wird hier zur Veranschaulichung verwendet, Sie können aber auch andere Formate wie CSV verwenden. In den folgenden Zellen sehen Sie, wie Sie mehrere Hyperspace-Indizes für dieses Beispieldataset erstellen und Spark anweisen, diese bei der Ausführung von Abfragen zu verwenden.
Die Beispieldatensätze entsprechen zwei Datasets: Abteilung und Mitarbeiter. Sie sollten die Pfade „emp_Location“ und „dept_Location“ so konfigurieren, dass sie für das Speicherkonto auf den gewünschten Speicherort zum Speichern der generierten Datendateien verweisen.
Die Ausgabe der Ausführung der folgenden Zelle zeigt den Inhalt der Datasets als Listen von Dreiergruppen an, gefolgt von Verweisen auf Datenrahmen, die erstellt wurden, um den Inhalt der einzelnen Datasets am bevorzugten Speicherort zu speichern.
import org.apache.spark.sql.DataFrame
// Sample department records
val departments = Seq(
(10, "Accounting", "New York"),
(20, "Research", "Dallas"),
(30, "Sales", "Chicago"),
(40, "Operations", "Boston"))
// Sample employee records
val employees = Seq(
(7369, "SMITH", 20),
(7499, "ALLEN", 30),
(7521, "WARD", 30),
(7566, "JONES", 20),
(7698, "BLAKE", 30),
(7782, "CLARK", 10),
(7788, "SCOTT", 20),
(7839, "KING", 10),
(7844, "TURNER", 30),
(7876, "ADAMS", 20),
(7900, "JAMES", 30),
(7934, "MILLER", 10),
(7902, "FORD", 20),
(7654, "MARTIN", 30))
// Save sample data in the Parquet format
import spark.implicits._
val empData: DataFrame = employees.toDF("empId", "empName", "deptId")
val deptData: DataFrame = departments.toDF("deptId", "deptName", "location")
val emp_Location: String = "/<yourpath>/employees.parquet" //TODO ** customize this location path **
val dept_Location: String = "/<yourpath>/departments.parquet" //TODO ** customize this location path **
empData.write.mode("overwrite").parquet(emp_Location)
deptData.write.mode("overwrite").parquet(dept_Location)
from pyspark.sql.types import StructField, StructType, StringType, IntegerType
# Sample department records
departments = [(10, "Accounting", "New York"), (20, "Research", "Dallas"), (30, "Sales", "Chicago"), (40, "Operations", "Boston")]
# Sample employee records
employees = [(7369, "SMITH", 20), (7499, "ALLEN", 30), (7521, "WARD", 30), (7566, "JONES", 20), (7698, "BLAKE", 30)]
# Create a schema for the dataframe
dept_schema = StructType([StructField('deptId', IntegerType(), True), StructField('deptName', StringType(), True), StructField('location', StringType(), True)])
emp_schema = StructType([StructField('empId', IntegerType(), True), StructField('empName', StringType(), True), StructField('deptId', IntegerType(), True)])
departments_df = spark.createDataFrame(departments, dept_schema)
employees_df = spark.createDataFrame(employees, emp_schema)
#TODO ** customize this location path **
emp_Location = "/<yourpath>/employees.parquet"
dept_Location = "/<yourpath>/departments.parquet"
employees_df.write.mode("overwrite").parquet(emp_Location)
departments_df.write.mode("overwrite").parquet(dept_Location)
using Microsoft.Spark.Sql.Types;
// Sample department records
var departments = new List<GenericRow>()
{
new GenericRow(new object[] {10, "Accounting", "New York"}),
new GenericRow(new object[] {20, "Research", "Dallas"}),
new GenericRow(new object[] {30, "Sales", "Chicago"}),
new GenericRow(new object[] {40, "Operations", "Boston"})
};
// Sample employee records
var employees = new List<GenericRow>() {
new GenericRow(new object[] {7369, "SMITH", 20}),
new GenericRow(new object[] {7499, "ALLEN", 30}),
new GenericRow(new object[] {7521, "WARD", 30}),
new GenericRow(new object[] {7566, "JONES", 20}),
new GenericRow(new object[] {7698, "BLAKE", 30}),
new GenericRow(new object[] {7782, "CLARK", 10}),
new GenericRow(new object[] {7788, "SCOTT", 20}),
new GenericRow(new object[] {7839, "KING", 10}),
new GenericRow(new object[] {7844, "TURNER", 30}),
new GenericRow(new object[] {7876, "ADAMS", 20}),
new GenericRow(new object[] {7900, "JAMES", 30}),
new GenericRow(new object[] {7934, "MILLER", 10}),
new GenericRow(new object[] {7902, "FORD", 20}),
new GenericRow(new object[] {7654, "MARTIN", 30})
};
// Save sample data in the Parquet format
var departmentSchema = new StructType(new List<StructField>()
{
new StructField("deptId", new IntegerType()),
new StructField("deptName", new StringType()),
new StructField("location", new StringType())
});
var employeeSchema = new StructType(new List<StructField>()
{
new StructField("empId", new IntegerType()),
new StructField("empName", new StringType()),
new StructField("deptId", new IntegerType())
});
DataFrame empData = spark.CreateDataFrame(employees, employeeSchema);
DataFrame deptData = spark.CreateDataFrame(departments, departmentSchema);
string emp_Location = "/<yourpath>/employees.parquet"; //TODO ** customize this location path **
string dept_Location = "/<yourpath>/departments.parquet"; //TODO ** customize this location path **
empData.Write().Mode("overwrite").Parquet(emp_Location);
deptData.Write().Mode("overwrite").Parquet(dept_Location);
Ergebnis:
departments: Seq[(Int, String, String)] = List((10,Accounting,New York), (20,Research,Dallas), (30,Sales,Chicago), (40,Operations,Boston))
employees: Seq[(Int, String, Int)] = List((7369,SMITH,20), (7499,ALLEN,30), (7521,WARD,30), (7566,JONES,20), (7698,BLAKE,30), (7782,CLARK,10), (7788,SCOTT,20), (7839,KING,10), (7844,TURNER,30), (7876,ADAMS,20), (7900,JAMES,30), (7934,MILLER,10), (7902,FORD,20), (7654,MARTIN,30))
empData: org.apache.spark.sql.DataFrame = [empId: int, empName: string ... 1 more field]
deptData: org.apache.spark.sql.DataFrame = [deptId: int, deptName: string ... 1 more field]
emp_Location: String = /your-path/employees.parquet
dept_Location: String = /your-path/departments.parquet
Überprüfen Sie die Inhalte der erstellten Parquet-Dateien, um sicherzustellen, dass sie die erwarteten Datensätze im richtigen Format enthalten. Sie verwenden diese Datendateien später zum Erstellen von Hyperspace-Indizes und zum Ausführen von Beispielabfragen.
Beim Ausführen der folgenden Zelle wird eine Ausgabe erzeugt, die die Zeilen der Mitarbeiter- und Abteilungs-DataFrames in tabellarischer Form anzeigt. Es sollte 14 Mitarbeiter und 4 Abteilungen geben, die jeweils mit einer der Dreiergruppen übereinstimmen, die Sie in der vorherigen Zelle erstellt haben.
// emp_Location and dept_Location are the user defined locations above to save parquet files
val empDF: DataFrame = spark.read.parquet(emp_Location)
val deptDF: DataFrame = spark.read.parquet(dept_Location)
// Verify the data is available and correct
empDF.show()
deptDF.show()
# emp_Location and dept_Location are the user-defined locations above to save parquet files
emp_DF = spark.read.parquet(emp_Location)
dept_DF = spark.read.parquet(dept_Location)
# Verify the data is available and correct
emp_DF.show()
dept_DF.show()
// emp_Location and dept_Location are the user-defined locations above to save parquet files
DataFrame empDF = spark.Read().Parquet(emp_Location);
DataFrame deptDF = spark.Read().Parquet(dept_Location);
// Verify the data is available and correct
empDF.Show();
deptDF.Show();
Ergebnis:
empDF: org.apache.spark.sql.DataFrame = [empId: int, empName: string ... 1 more field]
deptDF: org.apache.spark.sql.DataFrame = [deptId: int, deptName: string ... 1 more field]
|EmpId|EmpName|DeptId|
|-----|-------|------|
| 7499| ALLEN| 30|
| 7521| WARD| 30|
| 7369| SMITH| 20|
| 7844| TURNER| 30|
| 7876| ADAMS| 20|
| 7900| JAMES| 30|
| 7934| MILLER| 10|
| 7839| KING| 10|
| 7566| JONES| 20|
| 7698| BLAKE| 30|
| 7782| CLARK| 10|
| 7788| SCOTT| 20|
| 7902| FORD| 20|
| 7654| MARTIN| 30|
|DeptId| DeptName|Location|
|------|----------|--------|
| 10|Accounting|New York|
| 40|Operations| Boston|
| 20| Research| Dallas|
| 30| Sales| Chicago|
Indizes
Mit Hyperspace können Sie Indizes für Datensätze erstellen, die aus persistierten Datendateien gescannt wurden. Nach erfolgreicher Erstellung wird den Hyperspace-Metadaten ein Eintrag hinzugefügt, der dem Index entspricht. Diese Metadaten werden später vom Optimierer von Apache Spark (mit unseren Erweiterungen) während der Abfrageverarbeitung verwendet, um die richtigen Indizes zu finden und zu verwenden.
Nachdem Indizes erstellt wurden, können Sie verschiedene Aktionen ausführen:
- Aktualisieren bei Änderungen an den zugrunde liegenden Daten: Sie können einen vorhandenen Index aktualisieren, um die Änderungen zu übernehmen.
- Löschen, wenn der Index nicht benötigt wird: Sie können einen vorläufigen Löschvorgang durchführen. Dabei wird der Index nicht physisch gelöscht, sondern als „gelöscht“ markiert, sodass er nicht mehr in Ihren Workloads verwendet wird.
- Bereinigen, wenn ein Index nicht mehr benötigt wird: Sie können einen Index bereinigen, was eine vollständige physische Löschung des Indexinhalts und der zugehörigen Metadaten aus den Hyperspace-Metadaten erzwingt.
Aktualisieren – Wenn sich die zugrunde liegenden Daten ändern, können Sie einen vorhandenen Index aktualisieren, um dies zu erfassen. Löschen: Wenn der Index nicht benötigt wird, können Sie einen vorläufigen Löschvorgang durchführen, d. h., der Index wird nicht physisch gelöscht, sondern als „gelöscht“ markiert, sodass er nicht mehr in Ihren Workloads verwendet wird.
Die folgenden Abschnitte zeigen, wie solche Vorgänge zur Indexverwaltung in Hyperspace durchgeführt werden können.
Zuerst müssen Sie die erforderlichen Bibliotheken importieren und eine Instanz von Hyperspace erstellen. Sie verwenden diese Instanz später zum Aufrufen verschiedener Hyperspace-APIs, um Indizes für Ihre Beispieldaten zu erstellen und diese Indizes zu ändern.
Die Ausgabe der Ausführung der folgenden Zelle zeigt einen Verweis auf die erstellte Instanz von Hyperspace.
// Create an instance of Hyperspace
import com.microsoft.hyperspace._
val hyperspace: Hyperspace = Hyperspace()
from hyperspace import *
# Create an instance of Hyperspace
hyperspace = Hyperspace(spark)
// Create an instance of Hyperspace
using Microsoft.Spark.Extensions.Hyperspace;
Hyperspace hyperspace = new Hyperspace(spark);
Ergebnis:
hyperspace: com.microsoft.hyperspace.Hyperspace = com.microsoft.hyperspace.Hyperspace@1432f740
Erstellen von Indizes
Sie müssen zwei Informationen angeben, um einen Hyperspace-Index zu erstellen:
- Einen Spark-Datenrahmen, der auf die zu indexierenden Daten verweist.
- Ein Indexkonfigurationsobjekt (IndexConfig), das den Indexnamen sowie die indizierten und eingeschlossenen Spalten des Index angibt.
Sie beginnen damit, drei Hyperspace-Indizes für die Beispieldaten zu erstellen: zwei Indizes für das Abteilungsdataset namens „deptIndex1“ und „deptIndex2“ und einen Index für das Mitarbeiterdataset namens „empIndex“. Für jeden Index benötigen Sie eine entsprechende Indexkonfiguration (IndexConfig), um den Namen zusammen mit den Spaltenlisten für die indizierten und eingeschlossenen Spalten zu erfassen. Durch Ausführen der folgenden Zelle werden diese Indexkonfigurationen erstellt und in der Ausgabe aufgelistet.
Hinweis
Eine Indexspalte ist eine Spalte, die in Ihren Filtern oder Joinbedingungen angezeigt wird. Eine einbezogene Spalte ist eine Spalte, die in Ihrer Auswahl oder Ihrem Projekt enthalten ist.
Folgendes gilt z. B. für die nachfolgende Abfrage:
SELECT X
FROM T
WHERE Y = 2
Y kann eine Indexspalte und X eine eingeschlossene Spalte sein.
// Create index configurations
import com.microsoft.hyperspace.index.IndexConfig
val empIndexConfig: IndexConfig = IndexConfig("empIndex", Seq("deptId"), Seq("empName"))
val deptIndexConfig1: IndexConfig = IndexConfig("deptIndex1", Seq("deptId"), Seq("deptName"))
val deptIndexConfig2: IndexConfig = IndexConfig("deptIndex2", Seq("location"), Seq("deptName"))
# Create index configurations
emp_IndexConfig = IndexConfig("empIndex1", ["deptId"], ["empName"])
dept_IndexConfig1 = IndexConfig("deptIndex1", ["deptId"], ["deptName"])
dept_IndexConfig2 = IndexConfig("deptIndex2", ["location"], ["deptName"])
using Microsoft.Spark.Extensions.Hyperspace.Index;
var empIndexConfig = new IndexConfig("empIndex", new string[] {"deptId"}, new string[] {"empName"});
var deptIndexConfig1 = new IndexConfig("deptIndex1", new string[] {"deptId"}, new string[] {"deptName"});
var deptIndexConfig2 = new IndexConfig("deptIndex2", new string[] {"location"}, new string[] {"deptName"});
Ergebnis:
empIndexConfig: com.microsoft.hyperspace.index.IndexConfig = [indexName: empIndex; indexedColumns: deptid; includedColumns: empname]
deptIndexConfig1: com.microsoft.hyperspace.index.IndexConfig = [indexName: deptIndex1; indexedColumns: deptid; includedColumns: deptname]
deptIndexConfig2: com.microsoft.hyperspace.index.IndexConfig = [indexName: deptIndex2; indexedColumns: location; includedColumns: deptname]
Jetzt erstellen Sie drei Indizes mithilfe Ihrer Indexkonfigurationen. Zu diesem Zweck rufen Sie den Befehl „createIndex“ für unsere Hyperspace-Instanz auf. Dieser Befehl erfordert eine Indexkonfiguration und den Datenrahmen, der die zu indizierenden Zeilen enthält. Bei Ausführung der folgenden Zelle werden drei Indizes erstellt.
// Create indexes from configurations
import com.microsoft.hyperspace.index.Index
hyperspace.createIndex(empDF, empIndexConfig)
hyperspace.createIndex(deptDF, deptIndexConfig1)
hyperspace.createIndex(deptDF, deptIndexConfig2)
# Create indexes from configurations
hyperspace.createIndex(emp_DF, emp_IndexConfig)
hyperspace.createIndex(dept_DF, dept_IndexConfig1)
hyperspace.createIndex(dept_DF, dept_IndexConfig2)
// Create indexes from configurations
hyperspace.CreateIndex(empDF, empIndexConfig);
hyperspace.CreateIndex(deptDF, deptIndexConfig1);
hyperspace.CreateIndex(deptDF, deptIndexConfig2);
Indizes auflisten
Der folgende Code zeigt, wie Sie alle verfügbaren Indizes in einer Hyperspace-Instanz auflisten können. Dazu wird die API „Indizes“ verwendet, die Informationen über vorhandene Indizes als Spark-Datenrahmen zurückgibt, sodass Sie weitere Vorgänge ausführen können.
Sie können z. B. gültige Vorgänge für diesen Datenrahmen aufrufen, um seinen Inhalt zu prüfen oder ihn weiter zu analysieren (z. B. bestimmte Indizes filtern oder sie nach einer gewünschten Eigenschaft gruppieren).
Die folgende Zelle verwendet die show-Aktion des Datenrahmens, um die Zeilen vollständig auszugeben und Details der Indizes in tabellarischer Form anzuzeigen. Für jeden Index können Sie alle Informationen anzeigen, die Hyperspace über ihn in den Metadaten gespeichert hat. Sie bemerken sofort Folgendes:
- „config.indexName“, „config.indexedColumns“, „config.includedColumns“ und „status.status“ sind die Felder, auf die ein Benutzer normalerweise verweist.
- „dfSignature“ wird automatisch von Hyperspace generiert und ist für jeden Index eindeutig. Hyperspace verwendet diese Signatur intern, um den Index zu pflegen und ihn zur Abfragezeit auszunutzen.
In der folgenden Ausgabe sollten alle drei Indizes den Status „AKTIV“ aufweisen, und ihr Name, die indizierten Spalten und die eingeschlossenen Spalten sollten mit dem übereinstimmen, was Sie oben in den Indexkonfigurationen definiert haben.
hyperspace.indexes.show
hyperspace.indexes().show()
hyperspace.Indexes().Show();
Ergebnis:
|Config.IndexName|Config.IndexedColumns|Config.IncludedColumns| SchemaString| SignatureProvider| DfSignature| SerializedPlan|NumBuckets| DirPath|Status.Value|Stats.IndexSize|
|----------------|---------------------|----------------------|--------------------|--------------------|--------------------|--------------------|----------|--------------------|------------|---------------|
| deptIndex1| [deptId]| [deptName]|`deptId` INT,`dep...|com.microsoft.cha...|0effc1610ae2e7c49...|Relation[deptId#3...| 200|abfss://datasets@...| ACTIVE| 0|
| deptIndex2| [location]| [deptName]|`location` STRING...|com.microsoft.cha...|0effc1610ae2e7c49...|Relation[deptId#3...| 200|abfss://datasets@...| ACTIVE| 0|
| empIndex| [deptId]| [empName]|`deptId` INT,`emp...|com.microsoft.cha...|30768c6c9b2533004...|Relation[empId#32...| 200|abfss://datasets@...| ACTIVE| 0|
Indizes löschen
Sie können einen bestehenden Index löschen, indem Sie die API „deleteIndex“ verwenden und den Indexnamen angeben. Beim Löschen des Index wird ein vorläufiger Löschvorgang durchgeführt: Es wird hauptsächlich der Status des Index in den Hyperspace-Metadaten von „AKTIV“ auf „GELÖSCHT“ aktualisiert. Dadurch wird der gelöschte Index von jeder zukünftigen Abfrageoptimierung ausgeschlossen und Hyperspace wählt diesen Index nicht mehr für eine Abfrage aus.
Indexdateien für einen gelöschten Index bleiben jedoch weiterhin verfügbar (da es sich um ein vorläufiges Löschen handelt), sodass der Index auf Benutzerwunsch wiederhergestellt werden kann.
Die folgende Zelle löscht den Index „deptIndex2“ und listet danach Hyperspace-Metadaten auf. Die Ausgabe sollte der obigen Zelle für „Indizes auflisten“ ähnlich sein, mit Ausnahme von „deptIndex2“, dessen Status jetzt in „GELÖSCHT“ geändert werden sollte.
hyperspace.deleteIndex("deptIndex2")
hyperspace.indexes.show
hyperspace.deleteIndex("deptIndex2")
hyperspace.indexes().show()
hyperspace.DeleteIndex("deptIndex2");
hyperspace.Indexes().Show();
Ergebnis:
|Config.IndexName|Config.IndexedColumns|Config.IncludedColumns| SchemaString| SignatureProvider| DfSignature| SerializedPlan|NumBuckets| DirPath|Status.Value|Stats.IndexSize|
|----------------|---------------------|----------------------|--------------------|--------------------|--------------------|--------------------|----------|--------------------|------------|---------------|
| deptIndex1| [deptId]| [deptName]|`deptId` INT,`dep...|com.microsoft.cha...|0effc1610ae2e7c49...|Relation[deptId#3...| 200|abfss://datasets@...| ACTIVE| 0|
| deptIndex2| [location]| [deptName]|`location` STRING...|com.microsoft.cha...|0effc1610ae2e7c49...|Relation[deptId#3...| 200|abfss://datasets@...| DELETED| 0|
| empIndex| [deptId]| [empName]|`deptId` INT,`emp...|com.microsoft.cha...|30768c6c9b2533004...|Relation[empId#32...| 200|abfss://datasets@...| ACTIVE| 0|
Indizes wiederherstellen
Sie können die API „restoreIndex“ verwenden, um einen gelöschten Index wiederherzustellen. Dadurch wird die neueste Version des Index wieder in den Status „AKTIV“ versetzt und für Abfragen erneut verfügbar gemacht. Die folgende Zelle zeigt ein Beispiel für die Verwendung von „restoreIndex“. Sie löschen „deptIndex1“ und stellen ihn wieder her. Die Ausgabe zeigt, dass „deptIndex1“ nach dem Aufruf des Befehls „deleteIndex“ zunächst zum Status „GELÖSCHT“ wechselte und nach dem Aufruf von „restoreIndex“ zum Status „AKTIV“ zurückkehrte.
hyperspace.deleteIndex("deptIndex1")
hyperspace.indexes.show
hyperspace.restoreIndex("deptIndex1")
hyperspace.indexes.show
hyperspace.deleteIndex("deptIndex1")
hyperspace.indexes().show()
hyperspace.restoreIndex("deptIndex1")
hyperspace.indexes().show()
hyperspace.DeleteIndex("deptIndex1");
hyperspace.Indexes().Show();
hyperspace.RestoreIndex("deptIndex1");
hyperspace.Indexes().Show();
Ergebnis:
|Config.IndexName|Config.IndexedColumns|Config.IncludedColumns| SchemaString| SignatureProvider| DfSignature| SerializedPlan|NumBuckets| DirPath|Status.Value|Stats.indexSize|
|----------------|---------------------|----------------------|--------------------|--------------------|--------------------|--------------------|----------|--------------------|------------|---------------|
| deptIndex1| [deptId]| [deptName]|`deptId` INT,`dep...|com.microsoft.cha...|0effc1610ae2e7c49...|Relation[deptId#3...| 200|abfss://datasets@...| DELETED| 0|
| deptIndex2| [location]| [deptName]|`location` STRING...|com.microsoft.cha...|0effc1610ae2e7c49...|Relation[deptId#3...| 200|abfss://datasets@...| DELETED| 0|
| empIndex| [deptId]| [empName]|`deptId` INT,`emp...|com.microsoft.cha...|30768c6c9b2533004...|Relation[empId#32...| 200|abfss://datasets@...| ACTIVE| 0|
|Config.IndexName|Config.IndexedColumns|Config.IncludedColumns| SchemaString| SignatureProvider| DfSignature| SerializedPlan|NumBuckets| DirPath|Status.value|Stats.IndexSize|
|----------------|---------------------|----------------------|--------------------|--------------------|--------------------|--------------------|----------|--------------------|------------|---------------|
| deptIndex1| [deptId]| [deptName]|`deptId` INT,`dep...|com.microsoft.cha...|0effc1610ae2e7c49...|Relation[deptId#3...| 200|abfss://datasets@...| ACTIVE| 0|
| deptIndex2| [location]| [deptName]|`location` STRING...|com.microsoft.cha...|0effc1610ae2e7c49...|Relation[deptId#3...| 200|abfss://datasets@...| DELETED| 0|
| empIndex| [deptId]| [empName]|`deptId` INT,`emp...|com.microsoft.cha...|30768c6c9b2533004...|Relation[empId#32...| 200|abfss://datasets@...| ACTIVE| 0|
Indizes bereinigen
Sie können einen endgültigen Löschvorgang durchführen. Damit werden die Dateien und der Metadateneintrag für einen gelöschten Index mithilfe des Befehls vacuumIndex vollständig entfernt. Diese Aktion kann nicht rückgängig gemacht werden. Sie löscht alle Indexdateien physisch und wird deshalb als endgültiger Löschvorgang bezeichnet.
Die folgende Zelle bereinigt den Index „deptIndex2“ und zeigt nach dem Bereinigen Hyperspace-Metadaten an. Sie sollten Metadateneinträge für die beiden Indizes „deptIndex1“ und „empIndex“ sehen, beide mit dem Status „AKTIV“, und keinen Eintrag für „deptIndex2“.
hyperspace.vacuumIndex("deptIndex2")
hyperspace.indexes.show
hyperspace.vacuumIndex("deptIndex2")
hyperspace.indexes().show()
hyperspace.VacuumIndex("deptIndex2");
hyperspace.Indexes().Show();
Ergebnis:
|Config.IndexName|Config.IndexedColumns|Config.IncludedColumns| SchemaString| SignatureProvider| DfSignature| SerializedPlan|NumBuckets| DirPath|Status.Value|Stats.IndexSize|
|----------------|---------------------|----------------------|--------------------|--------------------|--------------------|--------------------|----------|--------------------|------------|---------------|
| deptIndex1| [deptId]| [deptName]|`deptId` INT,`dep...|com.microsoft.cha...|0effc1610ae2e7c49...|Relation[deptId#3...| 200|abfss://datasets@...| ACTIVE| 0|
| empIndex| [deptId]| [empName]|`deptId` INT,`emp...|com.microsoft.cha...|30768c6c9b2533004...|Relation[empId#32...| 200|abfss://datasets@...| ACTIVE| 0|
Aktivieren/deaktivieren von Hyperspace
Hyperspace bietet APIs zum Aktivieren oder Deaktivieren der Indexnutzung mit Spark.
- Mithilfe des Befehls enableHyperspace werden die Hyperspace-Optimierungsregeln für den Spark-Optimierer sichtbar, und sie nutzen vorhandene Hyperspace-Indizes zur Optimierung der Benutzerabfragen.
- Bei Verwendung des Befehls disableHyperspace gelten die Hyperspace-Regeln während der Abfrageoptimierung nicht mehr. Das Deaktivieren von Hyperspace hat keine Auswirkungen auf erstellte Indizes, da diese unverändert bleiben.
Die folgende Zelle zeigt, wie Sie Hyperspace mithilfe dieser Befehle aktivieren oder deaktivieren können. Die Ausgabe zeigt einen Verweis auf die vorhandene Spark-Sitzung, deren Konfiguration aktualisiert wird.
// Enable Hyperspace
spark.enableHyperspace
// Disable Hyperspace
spark.disableHyperspace
# Enable Hyperspace
Hyperspace.enable(spark)
# Disable Hyperspace
Hyperspace.disable(spark)
// Enable Hyperspace
spark.EnableHyperspace();
// Disable Hyperspace
spark.DisableHyperspace();
Ergebnis:
res48: org.apache.spark.sql.Spark™Session = org.apache.spark.sql.SparkSession@39fe1ddb
res51: org.apache.spark.sql.Spark™Session = org.apache.spark.sql.SparkSession@39fe1ddb
Indexnutzung
Damit Spark während der Abfrageverarbeitung Hyperspace-Indizes verwenden kann, müssen Sie sicherstellen, dass Hyperspace aktiviert ist.
Die folgende Zelle aktiviert Hyperspace und erstellt zwei Datenrahmen mit Ihren Beispieldatensätzen, die Sie zum Ausführen von Beispielabfragen verwenden können. Für jeden Datenrahmen werden einige Beispielzeilen ausgegeben.
// Enable Hyperspace
spark.enableHyperspace
val empDFrame: DataFrame = spark.read.parquet(emp_Location)
val deptDFrame: DataFrame = spark.read.parquet(dept_Location)
empDFrame.show(5)
deptDFrame.show(5)
# Enable Hyperspace
Hyperspace.enable(spark)
emp_DF = spark.read.parquet(emp_Location)
dept_DF = spark.read.parquet(dept_Location)
emp_DF.show(5)
dept_DF.show(5)
// Enable Hyperspace
spark.EnableHyperspace();
DataFrame empDFrame = spark.Read().Parquet(emp_Location);
DataFrame deptDFrame = spark.Read().Parquet(dept_Location);
empDFrame.Show(5);
deptDFrame.Show(5);
Ergebnis:
res53: org.apache.spark.sql.Spark™Session = org.apache.spark.sql.Spark™Session@39fe1ddb
empDFrame: org.apache.spark.sql.DataFrame = [empId: int, empName: string ... 1 more field]
deptDFrame: org.apache.spark.sql.DataFrame = [deptId: int, deptName: string ... 1 more field]
|empId|empName|deptId|
|-----|-------|------|
| 7499| ALLEN| 30|
| 7521| WARD| 30|
| 7369| SMITH| 20|
| 7844| TURNER| 30|
| 7876| ADAMS| 20|
Hier werden nur die oberen fünf Zeilen angezeigt.
|deptId| deptName|location|
|------|----------|--------|
| 10|Accounting|New York|
| 40|Operations| Boston|
| 20| Research| Dallas|
| 30| Sales| Chicago|
Indextypen
Gegenwärtig verfügt Hyperspace über Regeln zur Nutzung von Indizes für zwei Gruppen von Abfragen:
- Auswahlabfragen mit Filterprädikaten für Lookup oder Bereichsauswahl.
- Verknüpfen Sie Abfragen mit einem Gleichheits-Joinprädikat (d. h. Equijoins).
Indizes zur Beschleunigung von Filtern
Die erste Beispielabfrage führt einen Lookup für Datensätze zur Abteilung durch, wie in der folgenden Zelle gezeigt. In SQL sieht die Abfrage folgendermaßen aus:
SELECT deptName
FROM departments
WHERE deptId = 20
Das Ergebnis der Ausführung der folgenden Zelle zeigt:
- Abfrageergebnis, bei dem es sich um einen einzelnen Abteilungsnamen handelt.
- Abfrageplan, den Spark zur Ausführung der Abfrage verwendet hat.
Im Abfrageplan zeigt der FileScan-Operator am unteren Rand des Plans die Datenquelle an, aus der die Datensätze gelesen wurden. Der Speicherort dieser Datei gibt den Pfad zur neuesten Version des Index „deptIndex1“ an. Diese Informationen zeigen, dass Spark gemäß der Abfrage und mithilfe der Hyperspace-Optimierungsregeln zur Laufzeit beschlossen hat, den richtigen Index zu verwerten.
// Filter with equality predicate
val eqFilter: DataFrame = deptDFrame.filter("deptId = 20").select("deptName")
eqFilter.show()
eqFilter.explain(true)
# Filter with equality predicate
eqFilter = dept_DF.filter("""deptId = 20""").select("""deptName""")
eqFilter.show()
eqFilter.explain(True)
DataFrame eqFilter = deptDFrame.Filter("deptId = 20").Select("deptName");
eqFilter.Show();
eqFilter.Explain(true);
Ergebnis:
eqFilter: org.apache.spark.sql.DataFrame = [deptName: string]
|DeptName|
|--------|
|Research|
== Parsed Logical Plan ==
'Project [unresolvedalias('deptName, None)]
+- Filter (deptId#533 = 20)
+- Relation[deptId#533,deptName#534,location#535] parquet
== Analyzed Logical Plan ==
deptName: string
Project [deptName#534]
+- Filter (deptId#533 = 20)
+- Relation[deptId#533,deptName#534,location#535] parquet
== Optimized Logical Plan ==
Project [deptName#534]
+- Filter (isnotnull(deptId#533) && (deptId#533 = 20))
+- Relation[deptId#533,deptName#534,location#535] parquet
== Physical Plan ==
*(1) Project [deptName#534]
+- *(1) Filter (isnotnull(deptId#533) && (deptId#533 = 20))
+- *(1) FileScan parquet [deptId#533,deptName#534] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspaceon..., PartitionFilters: [], PushedFilters: [IsNotNull(deptId), EqualTo(deptId,20)], ReadSchema: struct<deptId:int,deptName:string>
Das zweite Beispiel ist eine Bereichsauswahlabfrage für Datensätze vom Typ „Abteilung“. In SQL sieht die Abfrage folgendermaßen aus:
SELECT deptName
FROM departments
WHERE deptId > 20
Ähnlich wie im ersten Beispiel zeigt die Ausgabe der folgenden Zelle die Abfrageergebnisse (Namen von zwei Abteilungen) und den Abfrageplan. Der Speicherort der Datendatei im FileScan-Operator zeigt, dass „deptIndex1“ zur Ausführung der Abfrage verwendet wurde.
// Filter with range selection predicate
val rangeFilter: DataFrame = deptDFrame.filter("deptId > 20").select("deptName")
rangeFilter.show()
rangeFilter.explain(true)
# Filter with range selection predicate
rangeFilter = dept_DF.filter("""deptId > 20""").select("deptName")
rangeFilter.show()
rangeFilter.explain(True)
// Filter with range selection predicate
DataFrame rangeFilter = deptDFrame.Filter("deptId > 20").Select("deptName");
rangeFilter.Show();
rangeFilter.Explain(true);
Ergebnis:
rangeFilter: org.apache.spark.sql.DataFrame = [deptName: string]
| DeptName|
|----------|
|Operations|
| Sales|
== Parsed Logical Plan ==
'Project [unresolvedalias('deptName, None)]
+- Filter (deptId#533 > 20)
+- Relation[deptId#533,deptName#534,location#535] parquet
== Analyzed Logical Plan ==
deptName: string
Project [deptName#534]
+- Filter (deptId#533 > 20)
+- Relation[deptId#533,deptName#534,location#535] parquet
== Optimized Logical Plan ==
Project [deptName#534]
+- Filter (isnotnull(deptId#533) && (deptId#533 > 20))
+- Relation[deptId#533,deptName#534,location#535] parquet
== Physical Plan ==
*(1) Project [deptName#534]
+- *(1) Filter (isnotnull(deptId#533) && (deptId#533 > 20))
+- *(1) FileScan parquet [deptId#533,deptName#534] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspaceon..., PartitionFilters: [], PushedFilters: [IsNotNull(deptId), GreaterThan(deptId,20)], ReadSchema: struct<deptId:int,deptName:string>
Das dritte Beispiel ist eine Abfrage, die Abteilungs- und Mitarbeiterdatensätze über die Abteilungs-ID verknüpft. Die entsprechende SQL-Anweisung ist unten dargestellt:
SELECT employees.deptId, empName, departments.deptId, deptName
FROM employees, departments
WHERE employees.deptId = departments.deptId
Das Ergebnis der Ausführung der folgenden Zelle zeigt die Abfrageergebnisse, nämlich die Namen von 14 Mitarbeitern und die Abteilung, in der jeder einzelne Mitarbeiter arbeitet. Der Abfrageplan ist ebenfalls in der Ausgabe enthalten. Beachten Sie, wie die Dateispeicherorte für zwei FileScan-Operatoren zeigen, dass Spark die Indizes „empIndex“ und „deptIndex1“ verwendet hat, um die Abfrage auszuführen.
// Join
val eqJoin: DataFrame =
empDFrame.
join(deptDFrame, empDFrame("deptId") === deptDFrame("deptId")).
select(empDFrame("empName"), deptDFrame("deptName"))
eqJoin.show()
eqJoin.explain(true)
# Join
eqJoin = emp_DF.join(dept_DF, emp_DF.deptId == dept_DF.deptId).select(emp_DF.empName, dept_DF.deptName)
eqJoin.show()
eqJoin.explain(True)
// Join
DataFrame eqJoin =
empDFrame
.Join(deptDFrame, empDFrame.Col("deptId") == deptDFrame.Col("deptId"))
.Select(empDFrame.Col("empName"), deptDFrame.Col("deptName"));
eqJoin.Show();
eqJoin.Explain(true);
Ergebnis:
eqJoin: org.apache.spark.sql.DataFrame = [empName: string, deptName: string]
|empName| deptName|
|-------|----------|
| SMITH| Research|
| JONES| Research|
| FORD| Research|
| ADAMS| Research|
| SCOTT| Research|
| KING|Accounting|
| CLARK|Accounting|
| MILLER|Accounting|
| JAMES| Sales|
| BLAKE| Sales|
| MARTIN| Sales|
| ALLEN| Sales|
| WARD| Sales|
| TURNER| Sales|
== Parsed Logical Plan ==
Project [empName#528, deptName#534]
+- Join Inner, (deptId#529 = deptId#533)
:- Relation[empId#527,empName#528,deptId#529] parquet
+- Relation[deptId#533,deptName#534,location#535] parquet
== Analyzed Logical Plan ==
empName: string, deptName: string
Project [empName#528, deptName#534]
+- Join Inner, (deptId#529 = deptId#533)
:- Relation[empId#527,empName#528,deptId#529] parquet
+- Relation[deptId#533,deptName#534,location#535] parquet
== Optimized Logical Plan ==
Project [empName#528, deptName#534]
+- Join Inner, (deptId#529 = deptId#533)
:- Project [empName#528, deptId#529]
: +- Filter isnotnull(deptId#529)
: +- Relation[empName#528,deptId#529] parquet
+- Project [deptId#533, deptName#534]
+- Filter isnotnull(deptId#533)
+- Relation[deptId#533,deptName#534] parquet
== Physical Plan ==
*(3) Project [empName#528, deptName#534]
+- *(3) SortMergeJoin [deptId#529], [deptId#533], Inner
:- *(1) Project [empName#528, deptId#529]
: +- *(1) Filter isnotnull(deptId#529)
: +- *(1) FileScan parquet [deptId#529,empName#528] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspaceon..., PartitionFilters: [], PushedFilters: [IsNotNull(deptId)], ReadSchema: struct<deptId:int,empName:string>, SelectedBucketsCount: 200 out of 200
+- *(2) Project [deptId#533, deptName#534]
+- *(2) Filter isnotnull(deptId#533)
+- *(2) FileScan parquet [deptId#533,deptName#534] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspaceon..., PartitionFilters: [], PushedFilters: [IsNotNull(deptId)], ReadSchema: struct<deptId:int,deptName:string>, SelectedBucketsCount: 200 out of 200
Unterstützung für SQL-Semantik
Die Indexverwendung ist hinsichtlich der Verwendung der Datenrahmen-API oder von Spark SQL transparent. Das folgende Beispiel zeigt dasselbe Join-Beispiel wie zuvor, in SQL-Form, wobei gegebenenfalls die Verwendung von Indizes gezeigt wird.
empDFrame.createOrReplaceTempView("EMP")
deptDFrame.createOrReplaceTempView("DEPT")
val joinQuery = spark.sql("SELECT EMP.empName, DEPT.deptName FROM EMP, DEPT WHERE EMP.deptId = DEPT.deptId")
joinQuery.show()
joinQuery.explain(true)
from pyspark.sql import SparkSession
emp_DF.createOrReplaceTempView("EMP")
dept_DF.createOrReplaceTempView("DEPT")
joinQuery = spark.sql("SELECT EMP.empName, DEPT.deptName FROM EMP, DEPT WHERE EMP.deptId = DEPT.deptId")
joinQuery.show()
joinQuery.explain(True)
empDFrame.CreateOrReplaceTempView("EMP");
deptDFrame.CreateOrReplaceTempView("DEPT");
var joinQuery = spark.Sql("SELECT EMP.empName, DEPT.deptName FROM EMP, DEPT WHERE EMP.deptId = DEPT.deptId");
joinQuery.Show();
joinQuery.Explain(true);
Ergebnis:
joinQuery: org.apache.spark.sql.DataFrame = [empName: string, deptName: string]
|empName| deptName|
|-------|----------|
| SMITH| Research|
| JONES| Research|
| FORD| Research|
| ADAMS| Research|
| SCOTT| Research|
| KING|Accounting|
| CLARK|Accounting|
| MILLER|Accounting|
| JAMES| Sales|
| BLAKE| Sales|
| MARTIN| Sales|
| ALLEN| Sales|
| WARD| Sales|
| TURNER| Sales|
== Parsed Logical Plan ==
'Project ['EMP.empName, 'DEPT.deptName]
+- 'Filter ('EMP.deptId = 'DEPT.deptId)
+- 'Join Inner
:- 'UnresolvedRelation `EMP`
+- 'UnresolvedRelation `DEPT`
== Analyzed Logical Plan ==
empName: string, deptName: string
Project [empName#528, deptName#534]
+- Filter (deptId#529 = deptId#533)
+- Join Inner
:- SubqueryAlias `emp`
: +- Relation[empId#527,empName#528,deptId#529] parquet
+- SubqueryAlias `dept`
+- Relation[deptId#533,deptName#534,location#535] parquet
== Optimized Logical Plan ==
Project [empName#528, deptName#534]
+- Join Inner, (deptId#529 = deptId#533)
:- Project [empName#528, deptId#529]
: +- Filter isnotnull(deptId#529)
: +- Relation[empId#527,empName#528,deptId#529] parquet
+- Project [deptId#533, deptName#534]
+- Filter isnotnull(deptId#533)
+- Relation[deptId#533,deptName#534,location#535] parquet
== Physical Plan ==
*(5) Project [empName#528, deptName#534]
+- *(5) SortMergeJoin [deptId#529], [deptId#533], Inner
:- *(2) Sort [deptId#529 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(deptId#529, 200)
: +- *(1) Project [empName#528, deptId#529]
: +- *(1) Filter isnotnull(deptId#529)
: +- *(1) FileScan parquet [deptId#529,empName#528] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspaceon..., PartitionFilters: [], PushedFilters: [IsNotNull(deptId)], ReadSchema: struct<deptId:int,empName:string>
+- *(4) Sort [deptId#533 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(deptId#533, 200)
+- *(3) Project [deptId#533, deptName#534]
+- *(3) Filter isnotnull(deptId#533)
+- *(3) FileScan parquet [deptId#533,deptName#534] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/your-path/departments.parquet], PartitionFilters: [], PushedFilters: [IsNotNull(deptId)], ReadSchema: struct<deptId:int,deptName:string>
API erläutern
Indizes sind hervorragend, aber woher wissen Sie, ob sie verwendet werden? Hyperspace ermöglicht es Benutzern, ihren ursprünglichen Plan mit dem aktualisierten indexabhängigen Plan zu vergleichen, bevor sie ihre Abfrage ausführen. Sie haben die Möglichkeit, zur Anzeige der Befehlsausgabe zwischen den Modi HTML, Klartext oder Konsole zu wählen.
Die folgende Zelle zeigt ein Beispiel mit HTML. Der hervorgehobene Abschnitt stellt den Unterschied zwischen ursprünglichen und aktualisierten Plänen zusammen mit den verwendeten Indizes dar.
spark.conf.set("spark.hyperspace.explain.displayMode", "html")
hyperspace.explain(eqJoin)(displayHTML(_))
eqJoin = emp_DF.join(dept_DF, emp_DF.deptId == dept_DF.deptId).select(emp_DF.empName, dept_DF.deptName)
spark.conf.set("spark.hyperspace.explain.displayMode", "html")
hyperspace.explain(eqJoin, True, displayHTML)
spark.Conf().Set("spark.hyperspace.explain.displayMode", "html");
spark.Conf().Set("spark.hyperspace.explain.displayMode.highlight.beginTag", "<b style=\"background:LightGreen\">");
spark.Conf().Set("spark.hyperspace.explain.displayMode.highlight.endTag", "</b>");
hyperspace.Explain(eqJoin, false, input => DisplayHTML(input));
Ergebnis:
Mit Indizes planen
Project [empName#528, deptName#534]
+- SortMergeJoin [deptId#529], [deptId#533], Inner
:- *(1) Project [empName#528, deptId#529]
: +- *(1) Filter isnotnull(deptId#529)
: +- *(1) FileScan parquet [deptId#529,empName#528] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspaceon..., PartitionFilters: [], PushedFilters: [IsNotNull(deptId)], ReadSchema: struct
+- *(2) Project [deptId#533, deptName#534]
+- *(2) Filter isnotnull(deptId#533)
+- *(2) FileScan parquet [deptId#533,deptName#534] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspaceon..., PartitionFilters: [], PushedFilters: [IsNotNull(deptId)], ReadSchema: struct
Ohne Indizes planen
Project [empName#528, deptName#534]
+- SortMergeJoin [deptId#529], [deptId#533], Inner
:- *(2) Sort [deptId#529 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(deptId#529, 200)
: +- *(1) Project [empName#528, deptId#529]
: +- *(1) Filter isnotnull(deptId#529)
: +- *(1) FileScan parquet [empName#528,deptId#529] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/your-path/employees.parquet], PartitionFilters: [], PushedFilters: [IsNotNull(deptId)], ReadSchema: struct
+- *(4) Sort [deptId#533 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(deptId#533, 200)
+- *(3) Project [deptId#533, deptName#534]
+- *(3) Filter isnotnull(deptId#533)
+- *(3) FileScan parquet [deptId#533,deptName#534] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/your-path/departments.parquet], PartitionFilters: [], PushedFilters: [IsNotNull(deptId)], ReadSchema: struct
Verwendete Indizes
deptIndex1:abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/<container>/indexes/public/deptIndex1/v__=0
empIndex:abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/<container>/indexes/public/empIndex/v__=0
Indizes aktualisieren
Wenn sich die ursprünglichen Daten, für die ein Index erstellt wurde, ändern, erfasst der Index nicht mehr den neuesten Stand der Daten. Sie können einen solchen veralteten Index mithilfe des Befehls refreshIndex aktualisieren. Dieser Befehl bewirkt, dass der Index vollständig neu erstellt und entsprechend den neuesten Datensätzen aktualisiert wird. Wir zeigen Ihnen in anderen Notebooks, wie Sie Ihren Index inkrementell aktualisieren können.
Die beiden folgenden Zellen zeigen ein Beispiel für dieses Szenario:
- Die erste Zelle fügt den ursprünglichen Abteilungsdaten zwei weitere Abteilungen hinzu. Sie liest eine Liste der Abteilungen und gibt sie aus, um zu überprüfen, ob die neuen Abteilungen ordnungsgemäß hinzugefügt wurden. Die Ausgabe zeigt insgesamt sechs Abteilungen: vier alte und zwei neue. Der Aufruf von refreshIndex aktualisiert „deptIndex1“, sodass der Index neue Abteilungen erfasst.
- In der zweiten Zelle wird unser Beispiel einer Bereichsauswahlabfrage ausgeführt. Die Ergebnisse sollten jetzt vier Abteilungen enthalten: Zwei haben Sie bereits bei der obigen Abfrage gesehen, und zwei sind die neuen Abteilungen, die Sie gerade hinzugefügt haben.
Spezifische Indexaktualisierung
val extraDepartments = Seq(
(50, "Inovation", "Seattle"),
(60, "Human Resources", "San Francisco"))
val extraDeptData: DataFrame = extraDepartments.toDF("deptId", "deptName", "location")
extraDeptData.write.mode("Append").parquet(dept_Location)
val deptDFrameUpdated: DataFrame = spark.read.parquet(dept_Location)
deptDFrameUpdated.show(10)
hyperspace.refreshIndex("deptIndex1")
extra_Departments = [(50, "Inovation", "Seattle"), (60, "Human Resources", "San Francisco")]
extra_departments_df = spark.createDataFrame(extra_Departments, dept_schema)
extra_departments_df.write.mode("Append").parquet(dept_Location)
dept_DFrame_Updated = spark.read.parquet(dept_Location)
dept_DFrame_Updated.show(10)
var extraDepartments = new List<GenericRow>()
{
new GenericRow(new object[] {50, "Inovation", "Seattle"}),
new GenericRow(new object[] {60, "Human Resources", "San Francisco"})
};
DataFrame extraDeptData = spark.CreateDataFrame(extraDepartments, departmentSchema);
extraDeptData.Write().Mode("Append").Parquet(dept_Location);
DataFrame deptDFrameUpdated = spark.Read().Parquet(dept_Location);
deptDFrameUpdated.Show(10);
hyperspace.RefreshIndex("deptIndex1");
Ergebnis:
extraDepartments: Seq[(Int, String, String)] = List((50,Inovation,Seattle), (60,Human Resources,San Francisco))
extraDeptData: org.apache.spark.sql.DataFrame = [deptId: int, deptName: string ... 1 more field]
deptDFrameUpdated: org.apache.spark.sql.DataFrame = [deptId: int, deptName: string ... 1 more field]
|deptId| deptName| location|
|------|---------------|-------------|
| 60|Human Resources|San Francisco|
| 10| Accounting| New York|
| 50| Inovation| Seattle|
| 40| Operations| Boston|
| 20| Research| Dallas|
| 30| Sales| Chicago|
Bereichsauswahl
val newRangeFilter: DataFrame = deptDFrameUpdated.filter("deptId > 20").select("deptName")
newRangeFilter.show()
newRangeFilter.explain(true)
newRangeFilter = dept_DFrame_Updated.filter("deptId > 20").select("deptName")
newRangeFilter.show()
newRangeFilter.explain(True)
DataFrame newRangeFilter = deptDFrameUpdated.Filter("deptId > 20").Select("deptName");
newRangeFilter.Show();
newRangeFilter.Explain(true);
Ergebnis:
newRangeFilter: org.apache.spark.sql.DataFrame = [deptName: string]
| DeptName|
|---------------|
|Human Resources|
| Inovation|
| Operations|
| Sales|
== Parsed Logical Plan ==
'Project [unresolvedalias('deptName, None)]
+- Filter (deptId#674 > 20)
+- Relation[deptId#674,deptName#675,location#676] parquet
== Analyzed Logical Plan ==
deptName: string
Project [deptName#675]
+- Filter (deptId#674 > 20)
+- Relation[deptId#674,deptName#675,location#676] parquet
== Optimized Logical Plan ==
Project [deptName#675]
+- Filter (isnotnull(deptId#674) && (deptId#674 > 20))
+- Relation[deptId#674,deptName#675,location#676] parquet
== Physical Plan ==
*(1) Project [deptName#675]
+- *(1) Filter (isnotnull(deptId#674) && (deptId#674 > 20))
+- *(1) FileScan parquet [deptId#674,deptName#675] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspaceon..., PartitionFilters: [], PushedFilters: [IsNotNull(deptId), GreaterThan(deptId,20)], ReadSchema: struct<deptId:int,deptName:string>
Hybridüberprüfung für veränderliche Datasets
Häufig, wenn in Ihren zugrunde liegenden Quelldaten neue Dateien angefügt oder vorhandene Dateien gelöscht wurden, veraltet Ihr Index und Hyperspace beschließt, ihn nicht zu verwenden. Es gibt jedoch Zeiten, wenn Sie den Index lediglich verwenden möchten, ohne ihn jedes Mal aktualisieren zu müssen. Hierfür kann es mehrere Gründe geben:
- Sie möchten den Index nicht fortlaufend aktualisieren, sondern stattdessen regelmäßig, weil Sie Ihre Workloads am besten verstehen.
- Sie haben nur ein paar Dateien hinzugefügt/entfernt und möchten nicht warten, bis ein weiterer Aktualisierungsauftrag abgeschlossen ist.
Damit Sie weiterhin einen veralteten Index verwenden dürfen, führt Hyperspace eine Hybridüberprüfung ein, eine neue Methode, mit der Benutzer veraltete oder abgelaufene Indizes verwenden können (z. B., wenn in den zugrunde liegenden Quelldaten neue Dateien angefügt oder vorhandene Dateien gelöscht wurden), ohne die Indizes zu aktualisieren.
Hierzu ändert Hyperspace, wenn Sie die entsprechende Konfiguration zum Aktivieren der Hybridüberprüfung festlegen, den Abfrageplan so, dass die Änderungen wie folgt genutzt werden:
- Angefügte Dateien können zusammengeführt werden, um Daten mithilfe von Union oder BucketUnion (für Join) zu indizieren. Das Mischen angefügter Daten kann auch vor dem Zusammenführen angewendet werden, falls erforderlich.
- Gelöschte Dateien können durch Einfügen einer Filterbedingung „NOT IN“ in die Abstammungsspalte der Indexdaten verarbeitet werden, sodass die indizierten Zeilen aus den gelöschten Dateien während der Abfrage ausgeschlossen werden können.
Sie können die Transformation des Abfrageplans in den folgenden Beispielen überprüfen.
Hinweis
Aktuell wird die Hybridüberprüfung nur für nicht partitionierte Daten unterstützt.
Hybridüberprüfung für angefügte Dateien – nicht partitionierte Daten
Im folgenden Beispiel werden nicht partitionierte Daten verwendet. In diesem Beispiel wird davon ausgegangen, dass der Join-Index für die Abfrage verwendet werden kann und dass BucketUnion für angefügte Dateien eingeführt wurde.
val testData = Seq(
("orange", 3, "2020-10-01"),
("banana", 1, "2020-10-01"),
("carrot", 5, "2020-10-02"),
("beetroot", 12, "2020-10-02"),
("orange", 2, "2020-10-03"),
("banana", 11, "2020-10-03"),
("carrot", 3, "2020-10-03"),
("beetroot", 2, "2020-10-04"),
("cucumber", 7, "2020-10-05"),
("pepper", 20, "2020-10-06")
).toDF("name", "qty", "date")
val testDataLocation = s"$dataPath/productTable"
testData.write.mode("overwrite").parquet(testDataLocation)
val testDF = spark.read.parquet(testDataLocation)
testdata = [
("orange", 3, "2020-10-01"),
("banana", 1, "2020-10-01"),
("carrot", 5, "2020-10-02"),
("beetroot", 12, "2020-10-02"),
("orange", 2, "2020-10-03"),
("banana", 11, "2020-10-03"),
("carrot", 3, "2020-10-03"),
("beetroot", 2, "2020-10-04"),
("cucumber", 7, "2020-10-05"),
("pepper", 20, "2020-10-06")
]
testdata_location = data_path + "/productTable"
from pyspark.sql.types import StructField, StructType, StringType, IntegerType
testdata_schema = StructType([
StructField('name', StringType(), True),
StructField('qty', IntegerType(), True),
StructField('date', StringType(), True)])
test_df = spark.createDataFrame(testdata, testdata_schema)
test_df.write.mode("overwrite").parquet(testdata_location)
test_df = spark.read.parquet(testdata_location)
using Microsoft.Spark.Sql.Types;
var products = new List<GenericRow>() {
new GenericRow(new object[] {"orange", 3, "2020-10-01"}),
new GenericRow(new object[] {"banana", 1, "2020-10-01"}),
new GenericRow(new object[] {"carrot", 5, "2020-10-02"}),
new GenericRow(new object[] {"beetroot", 12, "2020-10-02"}),
new GenericRow(new object[] {"orange", 2, "2020-10-03"}),
new GenericRow(new object[] {"banana", 11, "2020-10-03"}),
new GenericRow(new object[] {"carrot", 3, "2020-10-03"}),
new GenericRow(new object[] {"beetroot", 2, "2020-10-04"}),
new GenericRow(new object[] {"cucumber", 7, "2020-10-05"}),
new GenericRow(new object[] {"pepper", 20, "2020-10-06"})
};
var productsSchema = new StructType(new List<StructField>()
{
new StructField("name", new StringType()),
new StructField("qty", new IntegerType()),
new StructField("date", new StringType())
});
DataFrame testData = spark.CreateDataFrame(products, productsSchema);
string testDataLocation = $"{dataPath}/productTable";
testData.Write().Mode("overwrite").Parquet(testDataLocation);
// CREATE INDEX
hyperspace.createIndex(testDF, IndexConfig("productIndex2", Seq("name"), Seq("date", "qty")))
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1)
val filter1 = testDF.filter("name = 'banana'")
val filter2 = testDF.filter("qty > 10")
val query = filter1.join(filter2, "name")
// Check Join index rule is applied properly.
hyperspace.explain(query)(displayHTML(_))
# CREATE INDEX
hyperspace.createIndex(test_df, IndexConfig("productIndex2", ["name"], ["date", "qty"]))
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1)
filter1 = test_df.filter("name = 'banana'")
filter2 = test_df.filter("qty > 10")
query = filter1.join(filter2, "name")
# Check Join index rule is applied properly.
hyperspace.explain(query, True, displayHTML)
// CREATE INDEX
DataFrame testDF = spark.Read().Parquet(testDataLocation);
var productIndex2Config = new IndexConfig("productIndex", new string[] {"name"}, new string[] {"date", "qty"});
hyperspace.CreateIndex(testDF, productIndex2Config);
// Check Join index rule is applied properly.
DataFrame filter1 = testDF.Filter("name = 'banana'");
DataFrame filter2 = testDF.Filter("qty > 10");
DataFrame query = filter1.Join(filter2, filter1.Col("name") == filter2.Col("name"));
query.Show();
hyperspace.Explain(query, true, input => DisplayHTML(input));
Ergebnis:
=============================================================
Plan with indexes:
=============================================================
Project [name#607, qty#608, date#609, qty#632, date#633]
+- SortMergeJoin [name#607], [name#631], Inner
:- *(1) Project [name#607, qty#608, date#609]
: +- *(1) Filter (isnotnull(name#607) && (name#607 = banana))
: +- *(1) FileScan parquet [name#607,date#609,qty#608] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/p..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct, SelectedBucketsCount: 1 out of 200
+- *(2) Project [name#631, qty#632, date#633]
+- *(2) Filter (((isnotnull(qty#632) && (qty#632 > 10)) && isnotnull(name#631)) && (name#631 = banana))
+- *(2) FileScan parquet [name#631,date#633,qty#632] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/p..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct, SelectedBucketsCount: 1 out of 200
=============================================================
Plan without indexes:
=============================================================
Project [name#607, qty#608, date#609, qty#632, date#633]
+- SortMergeJoin [name#607], [name#631], Inner
:- *(2) Sort [name#607 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(name#607, 200), [id=#453]
: +- *(1) Project [name#607, qty#608, date#609]
: +- *(1) Filter (isnotnull(name#607) && (name#607 = banana))
: +- *(1) FileScan parquet [name#607,qty#608,date#609] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
+- *(4) Sort [name#631 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(name#631, 200), [id=#459]
+- *(3) Project [name#631, qty#632, date#633]
+- *(3) Filter (((isnotnull(qty#632) && (qty#632 > 10)) && isnotnull(name#631)) && (name#631 = banana))
+- *(3) FileScan parquet [name#631,qty#632,date#633] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
=============================================================
Indexes used:
=============================================================
productIndex2:abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/productIndex2/v__=0
// Append new files.
val appendData = Seq(
("orange", 13, "2020-11-01"),
("banana", 5, "2020-11-01")).toDF("name", "qty", "date")
appendData.write.mode("append").parquet(testDataLocation)
# Append new files.
append_data = [
("orange", 13, "2020-11-01"),
("banana", 5, "2020-11-01")
]
append_df = spark.createDataFrame(append_data, testdata_schema)
append_df.write.mode("append").parquet(testdata_location)
// Append new files.
var appendProducts = new List<GenericRow>()
{
new GenericRow(new object[] {"orange", 13, "2020-11-01"}),
new GenericRow(new object[] {"banana", 5, "2020-11-01"})
};
DataFrame appendData = spark.CreateDataFrame(appendProducts, productsSchema);
appendData.Write().Mode("Append").Parquet(testDataLocation);
Die Hybridüberprüfung ist standardmäßig deaktiviert. Daher werden Sie feststellen, dass Hyperspace den Index nicht verwendet, da neue Daten angefügt wurden.
In der Ausgabe werden keine Planunterschiede angezeigt, daher gibt es keine Hervorhebung.
// Hybrid Scan configs are false by default.
spark.conf.set("spark.hyperspace.index.hybridscan.enabled", "false")
spark.conf.set("spark.hyperspace.index.hybridscan.delete.enabled", "false")
val testDFWithAppend = spark.read.parquet(testDataLocation)
val filter1 = testDFWithAppend.filter("name = 'banana'")
val filter2 = testDFWithAppend.filter("qty > 10")
val query = filter1.join(filter2, "name")
hyperspace.explain(query)(displayHTML(_))
query.show
# Hybrid Scan configs are false by default.
spark.conf.set("spark.hyperspace.index.hybridscan.enabled", "false")
spark.conf.set("spark.hyperspace.index.hybridscan.delete.enabled", "false")
test_df_with_append = spark.read.parquet(testdata_location)
filter1 = test_df_with_append.filter("name = 'banana'")
filter2 = test_df_with_append.filter("qty > 10")
query = filter1.join(filter2, "name")
hyperspace.explain(query, True, displayHTML)
query.show()
// Hybrid Scan configs are false by default.
spark.Conf().Set("spark.hyperspace.index.hybridscan.enabled", "false");
spark.Conf().Set("spark.hyperspace.index.hybridscan.delete.enabled", "false");
DataFrame testDFWithAppend = spark.Read().Parquet(testDataLocation);
DataFrame filter1 = testDFWithAppend.Filter("name = 'banana'");
DataFrame filter2 = testDFWithAppend.Filter("qty > 10");
DataFrame query = filter1.Join(filter2, filter1.Col("name") == filter2.Col("name"));
query.Show();
hyperspace.Explain(query, true, input => DisplayHTML(input));
Ergebnis:
=============================================================
Plan with indexes:
=============================================================
Project [name#678, qty#679, date#680, qty#685, date#686]
+- SortMergeJoin [name#678], [name#684], Inner
:- *(2) Sort [name#678 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(name#678, 200), [id=#589]
: +- *(1) Project [name#678, qty#679, date#680]
: +- *(1) Filter (isnotnull(name#678) && (name#678 = banana))
: +- *(1) FileScan parquet [name#678,qty#679,date#680] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
+- *(4) Sort [name#684 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(name#684, 200), [id=#595]
+- *(3) Project [name#684, qty#685, date#686]
+- *(3) Filter (((isnotnull(qty#685) && (qty#685 > 10)) && (name#684 = banana)) && isnotnull(name#684))
+- *(3) FileScan parquet [name#684,qty#685,date#686] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), EqualTo(name,banana), IsNotNull(name)], ReadSchema: struct
=============================================================
Plan without indexes:
=============================================================
Project [name#678, qty#679, date#680, qty#685, date#686]
+- SortMergeJoin [name#678], [name#684], Inner
:- *(2) Sort [name#678 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(name#678, 200), [id=#536]
: +- *(1) Project [name#678, qty#679, date#680]
: +- *(1) Filter (isnotnull(name#678) && (name#678 = banana))
: +- *(1) FileScan parquet [name#678,qty#679,date#680] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
+- *(4) Sort [name#684 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(name#684, 200), [id=#542]
+- *(3) Project [name#684, qty#685, date#686]
+- *(3) Filter (((isnotnull(qty#685) && (qty#685 > 10)) && (name#684 = banana)) && isnotnull(name#684))
+- *(3) FileScan parquet [name#684,qty#685,date#686] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), EqualTo(name,banana), IsNotNull(name)], ReadSchema: struct
+------+---+----------+---+----------+
| name|qty| date|qty| date|
+------+---+----------+---+----------+
|banana| 11|2020-10-03| 11|2020-10-03|
|banana| 5|2020-11-01| 11|2020-10-03|
|banana| 1|2020-10-01| 11|2020-10-03|
+------+---+----------+---+----------
Aktivieren der Hybridüberprüfung
Im Plan mit Indizes sehen Sie, dass die Exchange-Hashpartitionierung nur für angefügte Dateien erforderlich ist, sodass wir immer noch die „gemischten“ Indexdaten mit angefügten Dateien verwenden können. BucketUnion wird verwendet, um „gemischte“ angefügte Dateien mit den Indexdaten zusammenzuführen.
// Enable Hybrid Scan config. "delete" config is not necessary since we only appended data.
spark.conf.set("spark.hyperspace.index.hybridscan.enabled", "true")
spark.enableHyperspace
// Need to redefine query to recalculate the query plan.
val query = filter1.join(filter2, "name")
hyperspace.explain(query)(displayHTML(_))
query.show
# Enable Hybrid Scan config. "delete" config is not necessary.
spark.conf.set("spark.hyperspace.index.hybridscan.enabled", "true")
# Need to redefine query to recalculate the query plan.
query = filter1.join(filter2, "name")
hyperspace.explain(query, True, displayHTML)
query.show()
// Enable Hybrid Scan config. "delete" config is not necessary.
spark.Conf().Set("spark.hyperspace.index.hybridscan.enabled", "true");
spark.EnableHyperspace();
// Need to redefine query to recalculate the query plan.
DataFrame query = filter1.Join(filter2, filter1.Col("name") == filter2.Col("name"));
query.Show();
hyperspace.Explain(query, true, input => DisplayHTML(input));
Ergebnis:
=============================================================
Plan with indexes:
=============================================================
Project [name#678, qty#679, date#680, qty#732, date#733]
+- SortMergeJoin [name#678], [name#731], Inner
:- *(3) Sort [name#678 ASC NULLS FIRST], false, 0
: +- BucketUnion 200 buckets, bucket columns: [name]
: :- *(1) Project [name#678, qty#679, date#680]
: : +- *(1) Filter (isnotnull(name#678) && (name#678 = banana))
: : +- *(1) FileScan parquet [name#678,date#680,qty#679] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/p..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct, SelectedBucketsCount: 1 out of 200
: +- Exchange hashpartitioning(name#678, 200), [id=#775]
: +- *(2) Project [name#678, qty#679, date#680]
: +- *(2) Filter (isnotnull(name#678) && (name#678 = banana))
: +- *(2) FileScan parquet [name#678,date#680,qty#679] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
+- *(6) Sort [name#731 ASC NULLS FIRST], false, 0
+- BucketUnion 200 buckets, bucket columns: [name]
:- *(4) Project [name#731, qty#732, date#733]
: +- *(4) Filter (((isnotnull(qty#732) && (qty#732 > 10)) && isnotnull(name#731)) && (name#731 = banana))
: +- *(4) FileScan parquet [name#731,date#733,qty#732] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/p..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct, SelectedBucketsCount: 1 out of 200
+- Exchange hashpartitioning(name#731, 200), [id=#783]
+- *(5) Project [name#731, qty#732, date#733]
+- *(5) Filter (((isnotnull(qty#732) && (qty#732 > 10)) && isnotnull(name#731)) && (name#731 = banana))
+- *(5) FileScan parquet [name#731,date#733,qty#732] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
=============================================================
Plan without indexes:
=============================================================
Project [name#678, qty#679, date#680, qty#732, date#733]
+- SortMergeJoin [name#678], [name#731], Inner
:- *(2) Sort [name#678 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(name#678, 200), [id=#701]
: +- *(1) Project [name#678, qty#679, date#680]
: +- *(1) Filter (isnotnull(name#678) && (name#678 = banana))
: +- *(1) FileScan parquet [name#678,qty#679,date#680] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
+- *(4) Sort [name#731 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(name#731, 200), [id=#707]
+- *(3) Project [name#731, qty#732, date#733]
+- *(3) Filter (((isnotnull(qty#732) && (qty#732 > 10)) && isnotnull(name#731)) && (name#731 = banana))
+- *(3) FileScan parquet [name#731,qty#732,date#733] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
=============================================================
Indexes used:
=============================================================
productIndex2:abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/productIndex2/v__=0
+------+---+----------+---+----------+
| name|qty| date|qty| date|
+------+---+----------+---+----------+
|banana| 1|2020-10-01| 11|2020-10-03|
|banana| 11|2020-10-03| 11|2020-10-03|
|banana| 5|2020-11-01| 11|2020-10-03|
+------+---+----------+---+----------+
Inkrementelle Indexaktualisierung
Wenn Sie bereit sind, Ihre Indizes zu aktualisieren, aber nicht den gesamten Index neu erstellen möchten, unterstützt Hyperspace das Aktualisieren von Indizes auf inkrementelle Weise mithilfe der hs.refreshIndex("name", "incremental")-API. Dadurch entfällt die Notwendigkeit einer vollständigen Neuerstellung des Indexes, wobei zuvor erstellte Indexdateien genutzt und Indizes nur mit den neu hinzugefügten Daten aktualisiert werden.
Natürlich müssen Sie sicherstellen, dass Sie regelmäßig die ergänzende optimizeIndex-API verwenden (siehe unten), um sicherzustellen, dass es nicht zu Leistungsverschlechterungen kommt. Es wird empfohlen, die Optimierung mindestens einmal jedes 10. Mal aufzurufen, das Sie refreshIndex(..., "incremental") aufrufen, wobei angenommen wird, dass die von Ihnen hinzugefügten/entfernten Daten < 10 % des ursprünglichen Datasets sind. Wenn Ihr ursprünglicher Datensatz beispielsweise 100 GB groß ist und Sie Daten in Schritten von 1 GB hinzugefügt oder entfernt haben, können Sie refreshIndex 10 Mal aufrufen, bevor Sie optimizeIndex aufrufen. Beachten Sie, dass dieses Beispiel zur Veranschaulichung verwendet wird und Sie es für Ihre Workloads anpassen müssen.
Beachten Sie im folgenden Beispiel das Hinzufügen eines Sortierknotens im Abfrageplan, wenn Indizes verwendet werden. Dies liegt daran, dass für die angehängten Datendateien partielle Indizes erstellt werden, was dazu führt, dass Spark ein Sort einführt. Beachten Sie außerdem, dass Shuffle, d. h. Exchange, immer noch aus dem Plan entfernt ist, was Ihnen die entsprechende Beschleunigung bietet.
def query(): DataFrame = {
val testDFWithAppend = spark.read.parquet(testDataLocation)
val filter1 = testDFWithAppend.filter("name = 'banana'")
val filter2 = testDFWithAppend.filter("qty > 10")
filter1.join(filter2, "name")
}
hyperspace.refreshIndex("productIndex2", "incremental")
hyperspace.explain(query())(displayHTML(_))
query().show
def query():
test_df_with_append = spark.read.parquet(testdata_location)
filter1 = test_df_with_append.filter("name = 'banana'")
filter2 = test_df_with_append.filter("qty > 10")
return filter1.join(filter2, "name")
hyperspace.refreshIndex("productIndex2", "incremental")
hyperspace.explain(query(), True, displayHTML)
query().show()
Ergebnis:
=============================================================
Plan with indexes:
=============================================================
Project [name#820, qty#821, date#822, qty#827, date#828]
+- SortMergeJoin [name#820], [name#826], Inner
:- *(1) Sort [name#820 ASC NULLS FIRST], false, 0
: +- *(1) Project [name#820, qty#821, date#822]
: +- *(1) Filter (isnotnull(name#820) && (name#820 = banana))
: +- *(1) FileScan parquet [name#820,date#822,qty#821] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/p..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct, SelectedBucketsCount: 1 out of 200
+- *(2) Sort [name#826 ASC NULLS FIRST], false, 0
+- *(2) Project [name#826, qty#827, date#828]
+- *(2) Filter (((isnotnull(qty#827) && (qty#827 > 10)) && (name#826 = banana)) && isnotnull(name#826))
+- *(2) FileScan parquet [name#826,date#828,qty#827] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/p..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), EqualTo(name,banana), IsNotNull(name)], ReadSchema: struct, SelectedBucketsCount: 1 out of 200
=============================================================
Plan without indexes:
=============================================================
Project [name#820, qty#821, date#822, qty#827, date#828]
+- SortMergeJoin [name#820], [name#826], Inner
:- *(2) Sort [name#820 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(name#820, 200), [id=#927]
: +- *(1) Project [name#820, qty#821, date#822]
: +- *(1) Filter (isnotnull(name#820) && (name#820 = banana))
: +- *(1) FileScan parquet [name#820,qty#821,date#822] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
+- *(4) Sort [name#826 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(name#826, 200), [id=#933]
+- *(3) Project [name#826, qty#827, date#828]
+- *(3) Filter (((isnotnull(qty#827) && (qty#827 > 10)) && (name#826 = banana)) && isnotnull(name#826))
+- *(3) FileScan parquet [name#826,qty#827,date#828] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), EqualTo(name,banana), IsNotNull(name)], ReadSchema: struct
+------+---+----------+---+----------+
| name|qty| date|qty| date|
+------+---+----------+---+----------+
|banana| 1|2020-10-01| 11|2020-10-03|
|banana| 11|2020-10-03| 11|2020-10-03|
|banana| 5|2020-11-01| 11|2020-10-03|
+------+---+----------+---+----------+
Optimieren des Indexlayouts
Nach dem mehrfachen Aufrufen von inkrementellen Aktualisierungen für neu angefügten Daten (z. B. wenn der Benutzer in kleinen Batches in Daten schreibt oder in Streamingszenarien) tendiert die Anzahl der Indexdateien dazu, groß zu werden, was die Leistung des Index beeinträchtigt (Problem großer Anzahl von kleinen Dateien). Hyperspace bietet die hyperspace.optimizeIndex("indexName")-API, um das Indexlayout zu optimieren und das Problem mit großen Dateien zu verringern.
Beachten Sie im folgenden Plan, dass Hyperspace den zusätzlichen Sortierknoten im Abfrageplan entfernt hat. Durch die Optimierung kann das Sortieren für jeden Indexbucket vermieden werden, der nur eine Datei enthält. Dies ist jedoch nur dann der Fall, wenn ALLE Indexbuckets nach optimizeIndex höchstens eine Datei pro Bucket enthalten.
// Append some more data and call refresh again.
val appendData = Seq(
("orange", 13, "2020-11-01"),
("banana", 5, "2020-11-01")).toDF("name", "qty", "date")
appendData.write.mode("append").parquet(testDataLocation)
hyperspace.refreshIndex("productIndex2", "incremental")
# Append some more data and call refresh again.
append_data = [
("orange", 13, "2020-11-01"),
("banana", 5, "2020-11-01")
]
append_df = spark.createDataFrame(append_data, testdata_schema)
append_df.write.mode("append").parquet(testdata_location)
hyperspace.refreshIndex("productIndex2", "incremental"
// Call optimize. Ensure that Sort is removed after optimization (This is possible here because after optimize, in this case, every bucket contains only 1 file.).
hyperspace.optimizeIndex("productIndex2")
hyperspace.explain(query())(displayHTML(_))
# Call optimize. Ensure that Sort is removed after optimization (This is possible here because after optimize, in this case, every bucket contains only 1 file.).
hyperspace.optimizeIndex("productIndex2")
hyperspace.explain(query(), True, displayHTML)
Ergebnis:
=============================================================
Plan with indexes:
=============================================================
Project [name#954, qty#955, date#956, qty#961, date#962]
+- SortMergeJoin [name#954], [name#960], Inner
:- *(1) Project [name#954, qty#955, date#956]
: +- *(1) Filter (isnotnull(name#954) && (name#954 = banana))
: +- *(1) FileScan parquet [name#954,date#956,qty#955] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/p..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct, SelectedBucketsCount: 1 out of 200
+- *(2) Project [name#960, qty#961, date#962]
+- *(2) Filter (((isnotnull(qty#961) && (qty#961 > 10)) && isnotnull(name#960)) && (name#960 = banana))
+- *(2) FileScan parquet [name#960,date#962,qty#961] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/p..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct, SelectedBucketsCount: 1 out of 200
=============================================================
Plan without indexes:
=============================================================
Project [name#954, qty#955, date#956, qty#961, date#962]
+- SortMergeJoin [name#954], [name#960], Inner
:- *(2) Sort [name#954 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(name#954, 200), [id=#1070]
: +- *(1) Project [name#954, qty#955, date#956]
: +- *(1) Filter (isnotnull(name#954) && (name#954 = banana))
: +- *(1) FileScan parquet [name#954,qty#955,date#956] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
+- *(4) Sort [name#960 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(name#960, 200), [id=#1076]
+- *(3) Project [name#960, qty#961, date#962]
+- *(3) Filter (((isnotnull(qty#961) && (qty#961 > 10)) && isnotnull(name#960)) && (name#960 = banana))
+- *(3) FileScan parquet [name#960,qty#961,date#962] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
=============================================================
Indexes used:
=============================================================
productIndex2:abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/productIndex2/v__=3
Optimierungsmodi
Der Standardmodus für die Optimierung ist der „schnelle“ Modus, bei dem Dateien, die kleiner als ein vordefinierter Schwellenwert sind, zur Optimierung ausgewählt werden. Um die Wirkung der Optimierung zu maximieren, lässt Hyperspace einen weiteren Optimierungsmodus „vollständig“ zu, wie unten gezeigt. In diesem Modus werden ALLE Indexdateien für die Optimierung ausgewählt, unabhängig von ihrer Dateigröße, und es wird das bestmögliche Layout des Indexes erstellt. Dieser Modus ist auch langsamer als der Standardoptimierungsmodus, da hier mehr Daten verarbeitet werden.
hyperspace.optimizeIndex("productIndex2", "full")
hyperspace.explain(query())(displayHTML(_))
hyperspace.optimizeIndex("productIndex2", "full")
hyperspace.explain(query(), True, displayHTML)
Ergebnis:
=============================================================
Plan with indexes:
=============================================================
Project [name#1000, qty#1001, date#1002, qty#1007, date#1008]
+- SortMergeJoin [name#1000], [name#1006], Inner
:- *(1) Project [name#1000, qty#1001, date#1002]
: +- *(1) Filter (isnotnull(name#1000) && (name#1000 = banana))
: +- *(1) FileScan parquet [name#1000,date#1002,qty#1001] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/p..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct, SelectedBucketsCount: 1 out of 200
+- *(2) Project [name#1006, qty#1007, date#1008]
+- *(2) Filter (((isnotnull(qty#1007) && (qty#1007 > 10)) && isnotnull(name#1006)) && (name#1006 = banana))
+- *(2) FileScan parquet [name#1006,date#1008,qty#1007] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/p..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct, SelectedBucketsCount: 1 out of 200
=============================================================
Plan without indexes:
=============================================================
Project [name#1000, qty#1001, date#1002, qty#1007, date#1008]
+- SortMergeJoin [name#1000], [name#1006], Inner
:- *(2) Sort [name#1000 ASC NULLS FIRST], false, 0
: +- Exchange hashpartitioning(name#1000, 200), [id=#1160]
: +- *(1) Project [name#1000, qty#1001, date#1002]
: +- *(1) Filter (isnotnull(name#1000) && (name#1000 = banana))
: +- *(1) FileScan parquet [name#1000,qty#1001,date#1002] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
+- *(4) Sort [name#1006 ASC NULLS FIRST], false, 0
+- Exchange hashpartitioning(name#1006, 200), [id=#1166]
+- *(3) Project [name#1006, qty#1007, date#1008]
+- *(3) Filter (((isnotnull(qty#1007) && (qty#1007 > 10)) && isnotnull(name#1006)) && (name#1006 = banana))
+- *(3) FileScan parquet [name#1006,qty#1007,date#1008] Batched: true, Format: Parquet, Location: InMemoryFileIndex[abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/data-777519/prod..., PartitionFilters: [], PushedFilters: [IsNotNull(qty), GreaterThan(qty,10), IsNotNull(name), EqualTo(name,banana)], ReadSchema: struct
=============================================================
Indexes used:
=============================================================
productIndex2:abfss://datasets@hyperspacebenchmark.dfs.core.windows.net/hyperspace/indexes-777519/productIndex2/v__=4