Leistungsunterschied map () vs withColumn ()

Oct 21 2020

Ich habe eine Tabelle mit über 100 Spalten. Ich muss doppelte Anführungszeichen aus bestimmten Spalten entfernen. Ich habe zwei Möglichkeiten gefunden, dies mit withColumn () und map () zu tun.

Verwenden von withColumn ()

cols_to_fix = ["col1", ..., "col20"]
for col in cols_to_fix:
    df = df.withColumn(col, regexp_replace(df[col], "\"", ""))

Verwenden von map ()

def remove_quotes(row: Row) -> Row:
    row_as_dict = row.asDict()
    cols_to_fix = ["col1", ..., "col20"]
    for column in cols_to_fix:
        if row_as_dict[column]:
            row_as_dict[column] = re.sub("\"", "", str(row_as_dict[column]))
    return Row(**row_as_dict)
 
df = df.rdd.map(remove_quotes).toDF(df.schema)

Hier ist meine Frage. Ich habe festgestellt, dass die Verwendung von map () in einer Tabelle mit ~ 25 Millionen Datensätzen etwa viermal länger dauert als mit Column (). Ich werde es wirklich begrüßen, wenn ein anderer Stapelüberlaufbenutzer den Grund für den Leistungsunterschied erklären kann, damit ich in Zukunft ähnliche Fallstricke vermeiden kann.

Antworten

Colin Oct 21 2020 at 14:28

Zunächst ein Ratschlag: Konvertieren Sie DataFrame nicht in RDD und führen Sie einfach df.map (Ihre Funktion hier) aus. Dies kann viel Zeit sparen. die folgende Seitehttps://dzone.com/articles/apache-spark-3-reasons-why-you-should-not-use-rdds Dies würde uns viel Zeit sparen. Die wichtigste Schlussfolgerung ist, dass RDD bemerkenswert langsam ist als DataFrame / Dataset, ganz zu schweigen von der Zeit, die für die Konvertierung von DataFrame zu RDD verwendet wird.

Lassen Sie uns jetzt über map und withColumn sprechen, ohne dass DataFrame in RDD konvertiert werden muss. Fazit zuerst: Die Karte ist normalerweise 5x langsamer als mit Column. Der Grund dafür ist, dass die Kartenoperation immer eine Deserialisierung und Serialisierung umfasst, während withColumn eine interessierende Spalte bearbeiten kann. Um genau zu sein, sollte die Kartenoperation die Zeile in mehrere Teile deserialisieren, auf denen die Operation ausgeführt wird.

Ein Beispiel hier: Nehmen wir an, wir haben einen DataFrame, der wie folgt aussieht: + -------- + ----------- + | language | users_count | + -------- + ----------- + | Java | 20000 | | Python | 100000 | | Scala | 3000 | + -------- + ----------- + dann wollen wir alle Werte in der Spalte users_count um 1 erhöhen, wir können es so machen

    df.map(row => {
        val usersCount = row.getInt(1) + 1
        (row.getString(0), usersCount)
}).toDF("language", "users_count_incremented_by_1")

Im obigen Code müssen wir zuerst jede Zeile deserialisieren, um die Werte in der zweiten Spalte zu extrahieren. Danach geben wir die geänderten Werte aus und speichern sie als DataFrame (dieser Schritt erfordert die Serialisierung von (a, b) in Zeile (a, b) da DataFrame nichts anderes als ein DataSet of Rows ist). Weitere Informationen finden Sie im folgenden ausgezeichneten Artikelhttps://medium.com/@fqaiser94/udfs-vs-map-vs-custom-spark-native-functions-91ab2c154b44

map kann nicht mit der Spalte selbst arbeiten, sondern muss mit den Werten der Spalte arbeiten. Das Abrufen der Werte erfordert eine Deserialisierung und das Speichern als Datenrahmen erfordert eine Serialisierung.

Die Karte ist jedoch immer noch von großem Nutzen: Mit Hilfe der Kartenmethode können Benutzer sehr ausgefeilte Operationen implementieren, während nur integrierte Operationen ausgeführt werden können, wenn wir nur withColumn verwenden.

Zusammenfassend lässt sich sagen, dass die Karte langsamer, aber flexibler ist. Column ist sicherlich die effizienteste, während die Funktionalität eingeschränkt ist.