Errore: attributi risolti mancanti nel join
Sto usando pyspark per eseguire una joindelle due tabelle con una condizione di join relativamente complessa (utilizzando maggiore di / minore rispetto alle condizioni di join). Funziona bene, ma si interrompe non appena aggiungo un fillnacomando prima del join.
Il codice ha un aspetto simile a questo:
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')
)
Ciò si traduce in un errore come questo:
org.apache.spark.sql.AnalysisException: Attributi risolti col1 # 4765 mancanti da col1 # 6488 , col2 # 4766, col3 # 4768, colx # 4823, coly # 4830, colz # 4764 nell'operatore! Join LeftOuter, ( (( col1 # 4765 = colx # 4823) && (col2 # 4766 = coly # 4830)) && (col3 # 4768> = colz # 4764)). Gli attributi con lo stesso nome compaiono nell'operazione: col1. Verifica se vengono utilizzati gli attributi corretti.
Sembra che la scintilla non riconosca più col1dopo aver eseguito il fillna. (L'errore non si verifica se commento ciò.) Il problema è che ho bisogno di quella dichiarazione. (E in generale ho semplificato molto questo esempio.)
Ho esaminato questa domanda , ma queste risposte non funzionano per me. In particolare, l'utilizzo di .alias('a')after the fillnanon funziona perché quindi spark non riconosce il anella condizione di join.
Qualcuno potrebbe:
- Spiegare esattamente perché sta accadendo e come posso evitarlo in futuro?
- Mi consigli su come risolverlo?
Grazie in anticipo per il vostro aiuto.
Risposte
Che cosa sta succedendo?
Per "sostituire" i valori vuoti, viene creato un nuovo dataframe che contiene nuove colonne. Queste nuove colonne hanno gli stessi nomi di quelle precedenti ma sono effettivamente oggetti Spark completamente nuovi. Nel codice Scala puoi vedere che le colonne "modificate" sono quelle appena create mentre le colonne originali vengono rilasciate .
Un modo per vedere questo effetto è chiamare la spiegazione sul dataframe prima e dopo aver sostituito i valori vuoti:
df_a.explain()
stampe
== 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]
mentre
df_a.fillna(42, subset=['col1']).explain()
stampe
== 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]
Entrambi i piani contengono una colonna chiamata col1, ma nel primo caso viene chiamata la rappresentazione interna col1#6Lmentre viene chiamata la seconda col1#27L.
Quando la condizione di join df_a.col1 == df_b.colxora è associata alla colonna, col1#6Lil join fallirà se solo la colonna col1#27Lfa parte della tabella di sinistra.
Come si risolve il problema?
Il modo più ovvio sarebbe spostare l'operazione `fillna` prima della definizione della condizione di join:df_a = df_a.fillna('NA', subset=['col1'])
join_cond = [
df_a.col1 == df_b.colx,
[...]
Se ciò non è possibile o desiderato, è possibile modificare la condizione di join. Invece di usare una colonna da dataframe ( df_a.col1) puoi usare una colonna che non è associata a nessun dataframe usando la funzione col . Questa colonna funziona solo in base al suo nome e quindi ignora quando la colonna viene sostituita nel dataframe:
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
]
Lo svantaggio di questo secondo approccio è che i nomi delle colonne in entrambe le tabelle devono essere univoci.