Spark SQL: busque un valor en varias columnas

Aug 28 2020

Tengo un conjunto de datos de estado como el siguiente:

Quiero seleccionar todas las filas de este conjunto de datos que tienen "FALLO" en cualquiera de estas 5 columnas de estado.

Por lo tanto, quiero que el resultado contenga solo ID 1, 2, 4 ya que tienen FALLO en una de las columnas de Estado.

Supongo que en SQL podemos hacer algo como a continuación:

SELECT * FROM status WHERE "FAILURE" IN (Status1, Status2, Status3, Status4, Status5);

En Spark, sé que puedo hacer un filtro comparando cada columna de estado con "FAILURE"

status.filter(s => {s.Status1.equals(FAILURE) || s.Status2.equals(FAILURE) ... and so on..})

Pero me gustaría saber si hay una forma más inteligente de hacer esto en Spark SQL.

¡Gracias por adelantado!

Respuestas

1 LeoC Aug 29 2020 at 00:34

En caso de que haya muchas columnas para examinar, considere una función recursiva que cortocircuita en la primera coincidencia, como se muestra a continuación:

val df = Seq(
  (1, "T", "F", "T", "F"),
  (2, "T", "T", "T", "T"),
  (3, "T", "T", "F", "T")
).toDF("id", "c1", "c2", "c3", "c4")

import org.apache.spark.sql.Column

def checkFor(elem: Column, cols: List[Column]): Column = cols match {
  case Nil =>
    lit(true)
  case h :: tail =>
    when(h === elem, lit(false)).otherwise(checkFor(elem, tail))
}

val cols = df.columns.filter(_.startsWith("c")).map(col).toList

df.where(checkFor(lit("F"), cols)).show

// +---+---+---+---+---+
// | id| c1| c2| c3| c4|
// +---+---+---+---+---+
// |  2|  T|  T|  T|  T|
// +---+---+---+---+---+
thebluephantom Aug 28 2020 at 22:41

Un ejemplo similar se puede modificar y filtrar en la nueva columna agregada. Eso se lo dejo a usted, aquí verificando ceros excluyendo la primera columna:

import org.apache.spark.sql.functions._
import spark.implicits._

val df = sc.parallelize(Seq(
    ("r1", 0.0, 0.0, 0.0, 0.0),
    ("r2", 6.4, 4.9, 6.3, 7.1),
    ("r3", 4.2, 0.0, 7.2, 8.4),
    ("r4", 1.0, 2.0, 0.0, 0.0)
)).toDF("ID", "a", "b", "c", "d")

val count_some_val = df.columns.tail.map(x => when(col(x) === 0.0, 1).otherwise(0)).reduce(_ + _)     

val df2 = df.withColumn("some_val_count", count_some_val)
df2.filter(col("some_val_count") > 0).show(false)

Afaik no es posible detenerse cuando el primer partido se encuentra fácilmente, pero recuerdo a una persona más inteligente que yo mostrándome este enfoque con lazy existe que creo que se detiene en el primer encuentro de un partido. Así entonces, pero con un enfoque diferente, que me gusta:

import org.apache.spark.sql.functions._
import spark.implicits._

val df = sc.parallelize(Seq(
    ("r1", 0.0, 0.0, 0.0, 0.0),
    ("r2", 6.0, 4.9, 6.3, 7.1),
    ("r3", 4.2, 0.0, 7.2, 8.4),
    ("r4", 1.0, 2.0, 0.0, 0.0)
)).toDF("ID", "a", "b", "c", "d")

df.map{r => (r.getString(0),r.toSeq.tail.exists(c => 
             c.asInstanceOf[Double]==0))}
  .toDF("ID","ones")
  .show() 
vaquarkhan Aug 29 2020 at 01:19
        scala> import org.apache.spark.sql.functions._
        import org.apache.spark.sql.functions._

        scala> import spark.implicits._
        import spark.implicits._

        scala> val df = Seq(
             |     ("Prop1", "SUCCESS", "SUCCESS", "SUCCESS", "FAILURE" ,"SUCCESS"),
             |     ("Prop2", "SUCCESS", "FAILURE", "SUCCESS", "FAILURE", "SUCCESS"),
             |     ("Prop3", "SUCCESS", "SUCCESS", "SUCCESS", "SUCCESS", "SUCCESS" ),
             |     ("Prop4", "SUCCESS", "FAILURE", "SUCCESS", "FAILURE", "SUCCESS"),
             |     ("Prop5", "SUCCESS", "SUCCESS", "SUCCESS", "SUCCESS","SUCCESS")
             |    ).toDF("Name", "Status1", "Status2", "Status3", "Status4","Status5")
        df: org.apache.spark.sql.DataFrame = [Name: string, Status1: string ... 4 more fields]


        scala> df.show
        +-----+-------+-------+-------+-------+-------+
        | Name|Status1|Status2|Status3|Status4|Status5|
        +-----+-------+-------+-------+-------+-------+
        |Prop1|SUCCESS|SUCCESS|SUCCESS|FAILURE|SUCCESS|
        |Prop2|SUCCESS|FAILURE|SUCCESS|FAILURE|SUCCESS|
        |Prop3|SUCCESS|SUCCESS|SUCCESS|SUCCESS|SUCCESS|
        |Prop4|SUCCESS|FAILURE|SUCCESS|FAILURE|SUCCESS|
        |Prop5|SUCCESS|SUCCESS|SUCCESS|SUCCESS|SUCCESS|
        +-----+-------+-------+-------+-------+-------+


        scala> df.where($"Name".isin("Prop1","Prop4") and $"Status1".isin("SUCCESS","FAILURE")).show
        +-----+-------+-------+-------+-------+-------+
        | Name|Status1|Status2|Status3|Status4|Status5|
        +-----+-------+-------+-------+-------+-------+
        |Prop1|SUCCESS|SUCCESS|SUCCESS|FAILURE|SUCCESS|
        |Prop4|SUCCESS|FAILURE|SUCCESS|FAILURE|SUCCESS|
        +-----+-------+-------+-------+-------+-------+