spark broadcast variabile Map che dà valore nullo
Sep 22 2020
Sto usando java8 con Spark v2.4.1.
Sto cercando di utilizzare la variabile Broadcast Mapper cercare utilizzando come mostrato di seguito:
Dati in ingresso:
+-----+-----+-----+
|code1|code2|code3|
+-----+-----+-----+
|1 |7 | 5 |
|2 |7 | 4 |
|3 |7 | 3 |
|4 |7 | 2 |
|5 |7 | 1 |
+-----+-----+-----+
Uscita prevista:
+-----+-----+-----+
|code1|code2|code3|
+-----+-----+-----+
|1 |7 |51 |
|2 |7 |41 |
|3 |7 |31 |
|4 |7 |21 |
|5 |7 |11 |
+-----+-----+-----+
Il mio codice attuale con diverse soluzioni che ho provato:
Map<Integer,Integer> lookup_map= new HashMap<>();
lookup_map.put(1,11);
lookup_map.put(2,21);
lookup_map.put(3,31);
lookup_map.put(4,41);
lookup_map.put(5,51);
JavaSparkContext javaSparkContext = JavaSparkContext.fromSparkContext(sparkSession.sparkContext());
Broadcast<Map<Integer,Integer>> lookup_mapBcVar = javaSparkContext.broadcast(lookup_map);
Dataset<Row> resultDs= dataDs
.withColumn("floor_code3", floor(col("code3")))
.withColumn("floor_code3_int", floor(col("code3")).cast(DataTypes.IntegerType))
.withColumn("map_code3", lit(((Map<Integer, Integer>)lookup_mapBcVar.getValue()).get(col("floor_code3_int"))))
.withColumn("five", lit(((Map<Integer, Integer>)lookup_mapBcVar.getValue()).get(5)))
.withColumn("five_lit", lit(((Map<Integer, Integer>)lookup_mapBcVar.getValue()).get(lit(5).cast(DataTypes.IntegerType))));
L'output del codice corrente utilizzando:
resultDs.printSchema();
resultDs.show();
root
|-- code1: integer (nullable = true)
|-- code2: integer (nullable = true)
|-- code3: double (nullable = true)
|-- floor_code3: long (nullable = true)
|-- floor_code3_int: integer (nullable = true)
|-- map_code3: null (nullable = true)
|-- five: integer (nullable = false)
|-- five_lit: null (nullable = true)
+-----+-----+-----+-----------+---------------+---------+----+--------+
|code1|code2|code3|floor_code3|floor_code3_int|map_code3|five|five_lit|
+-----+-----+-----+-----------+---------------+---------+----+--------+
| 1| 7| 5.0| 5| 5| null| 51| null|
| 2| 7| 4.0| 4| 4| null| 51| null|
| 3| 7| 3.0| 3| 3| null| 51| null|
| 4| 7| 2.0| 2| 2| null| 51| null|
| 5| 7| 1.0| 1| 1| null| 51| null|
+-----+-----+-----+-----------+---------------+---------+----+--------+
Per ricreare i dati di input:
List<String[]> stringAsList = new ArrayList<>();
stringAsList.add(new String[] { "1","7","5" });
stringAsList.add(new String[] { "2","7","4" });
stringAsList.add(new String[] { "3","7","3" });
stringAsList.add(new String[] { "4","7","2" });
stringAsList.add(new String[] { "5","7","1" });
JavaSparkContext sparkContext = new JavaSparkContext(sparkSession.sparkContext());
JavaRDD<Row> rowRDD = sparkContext.parallelize(stringAsList).map((String[] row) -> RowFactory.create(row));
StructType schema = DataTypes
.createStructType(new StructField[] {
DataTypes.createStructField("code1", DataTypes.StringType, false),
DataTypes.createStructField("code2", DataTypes.StringType, false),
DataTypes.createStructField("code3", DataTypes.StringType, false)
});
Dataset<Row> dataDf= sparkSession.sqlContext().createDataFrame(rowRDD, schema).toDF();
Dataset<Row> dataDs = dataDf
.withColumn("code1", col("code1").cast(DataTypes.IntegerType))
.withColumn("code2", col("code2").cast(DataTypes.IntegerType))
.withColumn("code3", col("code3").cast(DataTypes.IntegerType));
Cosa sto facendo di sbagliato qui?
Scala Notebook per lo stesso qui
Risposte
1 roamer Oct 01 2020 at 22:20
lit() return Tipo di colonna, ma map.get richiede il tipo int che puoi fare in questo modo
val df: DataFrame = spark.sparkContext.parallelize(Range(0, 10000), 4).toDF("sentiment")
val map = new util.HashMap[Int, Int]()
map.put(1, 1)
map.put(2, 2)
map.put(3, 3)
val bf: Broadcast[util.HashMap[Int, Int]] = spark.sparkContext.broadcast(map)
df.rdd.map(x => {
val num = x.getInt(0)
(num, bf.value.get(num))
}).toDF("key", "add_key").show(false)