Spark gleichzeitige JDBC-Datenlesevorgänge

Jan 27 2023
Haben Sie jemals den Prozess der Implementierung von Spark in Ihrem Projekt durchlaufen und dabei die optimale Anzahl an Shuffle-Partitionen, die Speicherzuweisung für Treiber- und Executor-Instanzen, die Anzahl der Executor-Kerne und all diese lustigen Dinge ermittelt, nur um Daten aus einer JDBC-Quelle wie dieser zu lesen? ? Was ist das Problem mit dem Code oben? Nun ja ... es ist sehr langsam :) Wenn Sie wie im obigen Beispiel eine SQL-Abfrage in Spark ausführen, verwenden Sie nur einen Thread (Sie können in der Spark-Benutzeroberfläche sehen, dass während der Ausführung des Codes nur eine Aufgabe ausgeführt wird) und den erstellten Datenrahmen Habe eine Partition. Die meisten von uns verwenden Spark in ETL/ELT-Pipelines zur Parallelverarbeitung, daher ist das Lesen von Daten aus einer JDBC-Quelle mit nur einem Thread wahrscheinlich nicht das, was wir wollen.

Haben Sie jemals den Prozess der Implementierung von Spark in Ihrem Projekt durchlaufen und dabei die optimale Anzahl an Shuffle-Partitionen, die Speicherzuweisung für Treiber- und Executor-Instanzen, die Anzahl der Executor-Kerne und all diese lustigen Dinge ermittelt, nur um Daten aus einer JDBC-Quelle wie dieser zu lesen? ?

jdbcDF = spark.read \
    .format("jdbc") \
    .option("driver", "org.postgresql.Driver") \
    .option("url", "jdbc:postgresql:dbserver") \
    .option("user", os.environ['user']) \
    .option("password", os.environ['pass']) \
    .option("query", query) \
    .load()

Was ist das Problem mit dem Code oben? Nun ja ... es ist sehr langsam :) Wenn Sie wie im obigen Beispiel eine SQL-Abfrage in Spark ausführen, verwenden Sie nur einen Thread (Sie können in der Spark-Benutzeroberfläche sehen, dass während der Ausführung des Codes nur eine Aufgabe ausgeführt wird) und den erstellten Datenrahmen Habe eine Partition.

Die meisten von uns verwenden Spark in ETL/ELT-Pipelines zur Parallelverarbeitung, daher ist das Lesen von Daten aus einer JDBC-Quelle mit nur einem Thread wahrscheinlich nicht das, was wir wollen. Ganz zu schweigen davon, dass wir den Datenrahmen nach der Ausführung der Abfrage neu partitionieren müssen, da die Daten jetzt in einer Partition gespeichert sind. Wie lesen wir also tatsächlich Daten gleichzeitig in Spark?

# pass the original query to dbtable as a subquery 
final_query = f'({query}) as q'

# depends on your usecase, how many executor cores you have in the cluster
# what do you plan to do with the dataframe downstream, data size etc
partitionsCount = 100

# I'm hardcoding the date values here, but I'll show below how to identify
# these values at run time if you don't know them upfront
min_date = '2020-01-01'
max_date = '2022-12-31'

jdbcDF = spark.read \
    .format("jdbc") \
    .option("driver", "org.postgresql.Driver") \
    .option("url", "jdbc:postgresql:dbserver") \
    .option("user", os.environ['user']) \
    .option("password", os.environ['pass']) \
    .option("dbtable", final_query) \
    .option("numPartitions", partitionsCount) \
    .option("partitionColumn", "date") \
    .option("lowerBound", f"{min_date}") \
    .option("upperBound", f"{max_date}") \
    .load()

  1. dbtable:
    Basierend auf der offiziellen Dokumentation von Spark können wir, wenn wir Daten gleichzeitig lesen und entsprechend unseren Anforderungen partitionieren möchten, die Abfrageoption nicht mehr verwenden und müssen stattdessen unsere SQL-Abfrage über die dbtable- Option weiterleiten. Der einzige Unterschied zwischen diesen beiden Optionen (außer dem Namen) besteht darin, dass die Abfrage jetzt in Klammern eingeschlossen werden muss und einen Alias ​​haben muss – im Grunde wird die ursprüngliche Abfrage als Unterabfrage übergeben.
  2. numPartitions:
    Die maximale Anzahl von Partitionen, die für die Parallelität beim Lesen und Schreiben von Tabellen verwendet werden können. Dies bestimmt auch die maximale Anzahl gleichzeitiger JDBC-Verbindungen. Wenn die Anzahl der zu schreibenden Partitionen diesen Grenzwert überschreitet, verringern wir sie auf diesen Grenzwert, indem wir vor dem Schreiben „coalesce(numPartitions)“ aufrufen. ( Quelle ) Ich habe meinen im EMR-Cluster auf 100 eingestellt, aber eine gute Vorgehensweise wäre, die Partitionsnummer zwischen dem 1- und 4-fachen der Anzahl der Kerne festzulegen. Es handelt sich dabei jedoch nicht um eine in Stein gemeißelte Zahl. Testen Sie daher immer mit verschiedenen Werten, um herauszufinden, was für Sie funktioniert.
  3. partitionColumn:
    Dies ist wahrscheinlich der wichtigste Parameter und derjenige, der den größten Unterschied in der Laufzeit Ihrer Abfragen ausmacht. Beachten Sie Folgendes:
    - Die betreffende Spalte muss numerisch, datums- oder zeitgestempelt sein
    . - Für die bestmögliche Leistung müssen die Spaltenwerte möglichst gleichmäßig verteilt sein und eine hohe Kardinalität aufweisen. Wir möchten datenverzerrte Spalten vermeiden – darauf werde ich weiter unten näher eingehen.
    – Wenn Sie mehrere Spalten mit den oben genannten Eigenschaften haben, wählen Sie die indizierte Spalte aus.
  4. Untergrenze, Obergrenze :
    Beachten Sie, dass „lowerBound“ und „upperBound“ nur zur Festlegung des Partitionsschritts verwendet werden, nicht zum Filtern der Zeilen in der Tabelle. Daher werden alle Zeilen in der Tabelle partitioniert und zurückgegeben. ( Quelle ) Unten sehen Sie, wie Spark jede Partition generiert. Aus diesem Grund möchten Sie keine Spalte mit niedriger Kardinalität als Partitionsspalte verwenden, aber ich gebe unten ein Beispiel für den Fall, dass es jetzt keinen Sinn ergibt.
  5. Foto von Luminousmen auf Luminousmen.com

Stellen Sie sich vor, dass die Spalte, die Sie für die Partitionierung verwenden möchten, nur die Werte 0 und 1 hat – ein extremes Beispiel, aber ein häufiges Szenario im Data Engineering, bei dem die booleschen Werte „True“ und „ False“ in numerische Werte umgewandelt werden.

Da es nur Nullen und Einsen gibt, verwenden Sie 0 als Untergrenze und 1 als Obergrenze. Basierend auf der Art und Weise, wie die Partitionen generiert werden – siehe Bild oben – werden alle Zeilen mit dem Wert 0 in die erste Partition (die WHERE-Klausel der ersten Partition ist die einzige Klausel, in der 0 nicht herausgefiltert wird) und alle Zeilen verschoben mit dem Wert 1 wird in die letzte Partition verschoben (die WHERE-Klausel der letzten Partition ist die einzige Klausel, in der 1 nicht herausgefiltert wird) .

Unabhängig davon, wie viele Partitionen Sie erstellen, werden also nur zwei Partitionen bestückt und nur zwei Executor-Kerne erledigen die gesamte Arbeit. Die anderen befinden sich im Ruhezustand, wenn kein anderer Job für sie geplant ist. Aus diesem Grund wird die Leistung Ihrer Abfrage beeinträchtigt.

Datenschiefe

Stellen wir uns vor, wir haben einen Datensatz mit 20 Millionen Zeilen und 30 Partitionen, wobei die Unter- und Obergrenzen 2020–01–01 und 2022–12–31 sind. Für den Zweck dieses Artikels verwende ich eine Spark-Sitzung mit 30 Executor-Kernen (was bedeutet, dass wir bis zu 30 Aufgaben gleichzeitig ausführen können). Beachten Sie, dass Spark pro Partition eine Aufgabe zuweist, sodass jede Partition von einem Executor-Kern verarbeitet wird.

Die Daten jedes Jahres werden in 10 Partitionen aufgeteilt. Da 1 Aufgabe einer Partition zugewiesen wird, ist jede den Jahren 2020 und 2021 zugewiesene Aufgabe für die Verarbeitung von 100.000 Zeilen verantwortlich, jede dem Jahr 2022 zugewiesene Aufgabe muss jedoch verarbeitet werden 1,8 Mio. Zeilen. Das ist ein Anstieg der Daten, die von einer Aufgabe verarbeitet werden müssen, um 1700 %.

Das Ergebnis ist, dass die ersten 20 Aufgaben die Verarbeitung ihrer zugewiesenen Partition in kürzester Zeit abschließen und danach die 20 den abgeschlossenen Aufgaben zugewiesenen Executor-Kerne bis zum letzten Job im Leerlauf bleiben (sofern es keinen anderen Job gibt, bei dem die inaktiven Executor-Kerne verwendet werden können). 10 Aufgaben beenden die Verarbeitung von 1700 % mehr Datensätzen. Dies ist eine Verschwendung von Ressourcen und erhöht die Betriebskosten, insbesondere wenn Sie Spark in einem EMR-Cluster, Glue Job oder Databricks verwenden.

Wie können wir das beheben? Die einfachste Lösung: Verwenden Sie eine andere Spalte zur Partitionierung. Wenn die Verwendung dieser Spalte jedoch von entscheidender Bedeutung ist, können wir statt einer zwei Abfragen und zwei Datenrahmen verwenden. Sie werden immer noch Probleme mit verzerrten Daten im Nachhinein haben (ein wiederkehrendes Problem in parallelen Rechensystemen), aber wir gehen ein Problem nach dem anderen an.

Die erste Abfrage lädt die ersten beiden Jahre (die 30 Executor-Kerne werden zwischen 2020 und 2021 aufgeteilt), dann lädt die zweite Abfrage das letzte Jahr (also werden jetzt alle 30 Executor-Kerne für 2022 verwendet). Aus diesem Grund wird eine bestimmte Aufgabe im Zusammenhang mit 2022 nun für die Verarbeitung von 600.000 Zeilen anstelle von 1,8 Millionen Zeilen verantwortlich sein. Das ist eine Reduzierung der Zeilen, die von einem einzelnen Executor-Kern verarbeitet werden müssen, um 66 % im Vergleich zur ursprünglichen Implementierung. Darüber hinaus verwenden wir jetzt alle Executor-Kerne aus unserer Spark-Sitzung, anstatt sie im Leerlauf zu lassen.

Was passiert, wenn ich keine Zahlen-, Datums- oder Zeitstempelspalte habe oder die Spalte verzerrte Daten enthält?

Kein Problem. Als Dateningenieur row_number()sollte die Verwendung eines CTE und einer Fensterfunktion zum Generieren einer numerischen Spalte durch ein Kinderspiel sein. Beim Generieren der neuen Spalte wird ein gewisser Mehraufwand entstehen, aber selbst mit diesem Mehraufwand wird die endgültige Abfrage viel schneller ausgeführt, indem die neu erstellte Spalte verwendet wird, anstatt zu versuchen, alle Daten in einer Partition mit einem Executor-Kern zu laden.

with cte as
(
    SELECT f1, f2, f3
    FROM table
)
select *, row_number() over (order by f1) as rn from cte

original_query = "SELECT f1, f2, f3 FROM table"

final_query = f'(with cte as ({original_query}) select *, row_number() over (order by f1) as rn from cte) as q'

Wir führen zwei Abfragen anstelle einer aus:
– Die erste Abfrage wird ausgeführt, um die Mindest- und Höchstwerte der Spalte zu ermitteln, die wir als PartitionSpalte
verwenden . – In der zweiten Abfrage übergeben wir die Ergebnisse der ersten Abfrage als Argumente für die Unter- und Obergrenze .

query_min_max = f"SELECT min(date), max(date) from ({query}) q"
final_query = f"({query}) as q"

df_min_max = spark.read \
    .format("jdbc") \
    .option("driver", "org.postgresql.Driver") \
    .option("url", "jdbc:postgresql:dbserver") \
    .option("user", os.environ['user']) \
    .option("password", os.environ['pass']) \
    .option("query", query_min_max) \
    .load()

min = df_min_max.first()["min"]
max = df_min_max.first()["max"]

jdbcDF = spark.read \
    .format("jdbc") \
    .option("driver", "org.postgresql.Driver") \
    .option("url", "jdbc:postgresql:dbserver") \
    .option("user", os.environ['user']) \
    .option("password", os.environ['pass']) \
    .option("dbtable", final_query) \
    .option("numPartitions", partitionsCount) \
    .option("partitionColumn", "date") \
    .option("lowerBound", f"{min}") \
    .option("upperBound", f"{max}") \
    .load()

  • Wenn Sie Spark in Ihrer Anwendung implementieren, versuchen Sie, es optimal zu nutzen und gleichzeitig Daten zu lesen. Wenn Sie eine SQL-Abfrage an spark.read() ausführen, wie in der JDBC-Dokumentation von Spark gezeigt, werden alle Daten in einer Partition gepusht und nur ein Executor-Kern verwendet, unabhängig davon, wie viele Kerne Sie in Ihrer Spark-Sitzung eingerichtet haben.
  • Wenn Sie Spark heute verwenden, um Daten aus JDBC-Datenquellen mithilfe des Single-Threaded-Ansatzes zu lesen, versuchen Sie, Ihre Abfragen umzugestalten und ihre Leistung zu verbessern, erhalten Sie vielleicht sogar eine Beförderung, wenn Sie schon dabei sind :)
  • Letztendlich ist Zeit Geld. Wenn Sie sich Sorgen über den Aufwand machen, Ihre Abfragen neu zu schreiben, kann ich aus Erfahrung sprechen – einige der Spark-Abfragen, die ich mit den oben beschriebenen Techniken umgestaltet habe, werden jetzt um mehr als 90 % schneller ausgeführt als zuvor . Sie können sich vorstellen, wie positiv diese Verbesserung in meinem Team aufgenommen wurde.

Gabriel