Fehler: Behobene Attribute fehlen im Join
Ich verwende pyspark, um eine joinvon zwei Tabellen mit einer relativ komplexen Verknüpfungsbedingung auszuführen (wobei größer als / kleiner als in den Verknüpfungsbedingungen verwendet wird). Dies funktioniert einwandfrei, bricht jedoch zusammen, sobald ich fillnavor dem Join einen Befehl hinzufüge .
Der Code sieht ungefähr so aus:
join_cond = [
df_a.col1 == df_b.colx,
df_a.col2 == df_b.coly,
df_a.col3 >= df_b.colz
]
df = (
df_a
.fillna('NA', subset=['col1'])
.join(df_b, join_cond, 'left')
)
Dies führt zu einem Fehler wie dem folgenden:
org.apache.spark.sql.AnalysisException: Behobene Attribute col1 # 4765 fehlen in col1 # 6488 , col2 # 4766, col3 # 4768, colx # 4823, coly # 4830, colz # 4764 im Operator! Join LeftOuter, ( (( col1 # 4765 = colx # 4823) && (col2 # 4766 = coly # 4830)) && (col3 # 4768> = colz # 4764)). Gleichnamige Attribute erscheinen in der Operation: col1. Bitte überprüfen Sie, ob die richtigen Attribute verwendet werden.
Es sieht so aus, als würde der Funke col1nach dem Ausführen des nicht mehr erkannt fillna. (Der Fehler tritt nicht auf, wenn ich das auskommentiere.) Das Problem ist, dass ich diese Aussage brauche. (Und im Allgemeinen habe ich dieses Beispiel stark vereinfacht.)
Ich habe mir diese Frage angesehen , aber diese Antworten funktionieren bei mir nicht. Insbesondere funktioniert die Verwendung .alias('a')nach dem fillnanicht, da der Funke dann adie Verknüpfungsbedingung nicht erkennt .
Könnte jemand:
- Erklären Sie genau, warum dies geschieht und wie ich es in Zukunft vermeiden kann.
- Beraten Sie mich, wie ich es lösen kann?
Vielen Dank im Voraus für Ihre Hilfe.
Antworten
Was ist los?
Um leere Werte zu "ersetzen", wird ein neuer Datenrahmen erstellt, der neue Spalten enthält. Diese neuen Spalten haben dieselben Namen wie die alten, sind jedoch praktisch völlig neue Spark-Objekte. Im Scala-Code können Sie sehen, dass die "geänderten" Spalten neu erstellt wurden, während die ursprünglichen Spalten gelöscht werden .
Ein Weg , um diesen Effekt zu sehen , ist zu nennen erklärt auf dem Datenrahmen vor und nach den leeren Werten ersetzt:
df_a.explain()
druckt
== Physical Plan ==
*(1) Project [_1#0L AS col1#6L, _2#1L AS col2#7L, _3#2L AS col3#8L]
+- *(1) Scan ExistingRDD[_1#0L,_2#1L,_3#2L]
während
df_a.fillna(42, subset=['col1']).explain()
druckt
== Physical Plan ==
*(1) Project [coalesce(_1#0L, 42) AS col1#27L, _2#1L AS col2#7L, _3#2L AS col3#8L]
+- *(1) Scan ExistingRDD[_1#0L,_2#1L,_3#2L]
Beide Pläne enthalten eine Spalte mit dem Namen col1, aber im ersten Fall wird die interne Darstellung aufgerufen, col1#6Lwährend die zweite aufgerufen wird col1#27L.
Wenn die Verknüpfungsbedingung df_a.col1 == df_b.colxjetzt col1#6Lder Spalte zugeordnet col1#27List, schlägt die Verknüpfung fehl, wenn nur die Spalte Teil der linken Tabelle ist.
Wie kann das Problem gelöst werden?
Der naheliegende Weg wäre, die Operation "fillna" vor der Definition der Verknüpfungsbedingung zu verschieben:df_a = df_a.fillna('NA', subset=['col1'])
join_cond = [
df_a.col1 == df_b.colx,
[...]
Wenn dies nicht möglich oder gewünscht ist, können Sie die Join-Bedingung ändern. Anstatt eine Spalte aus dem Datenrahmen ( df_a.col1) zu verwenden, können Sie mithilfe der Funktion col eine Spalte verwenden, die keinem Datenrahmen zugeordnet ist . Diese Spalte funktioniert nur basierend auf ihrem Namen und wird daher ignoriert, wenn die Spalte im Datenrahmen ersetzt wird:
from pyspark.sql import functions as F
join_cond = [
F.col("col1") == df_b.colx,
df_a.col2 == df_b.coly,
df_a.col3 >= df_b.colz
]
Der Nachteil dieses zweiten Ansatzes ist, dass die Spaltennamen in beiden Tabellen eindeutig sein müssen.