【发布时间】:2020-12-10 18:46:41
【问题描述】:
在我之前的问题(Spark Structured Streaming dynamic lookup with Redis)中,感谢https://stackoverflow.com/users/689676/fe2s,我成功通过mapparttions到达redis
我尝试使用 mappartitions,但我无法解决一个问题,即如何在迭代时达到以下代码部分中的每一行列。 因为我想根据我保存在 Redis 中的查找字段来丰富我的每行。 我发现了类似的东西,但是我如何才能到达数据框列并添加查找 Redis 的新列。 非常感谢任何帮助,谢谢。
import org.apache.spark.sql.types._
def transformRow(row: Row): Row = {
Row.fromSeq(row.toSeq ++ Array[Any]("val1", "val2"))
}
def transformRows(iter: Iterator[Row]): Iterator[Row] =
{
val redisConn =new RedisClient("xxx.xxx.xx.xxx",6379,1,Option("Secret123"))
println(redisConn.get("ModelValidityPeriodName").getOrElse(""))
//want to reach DataFrame column here
redisConn.close()
iter.map(transformRow)
}
val newSchema = StructType(raw_customer_df.schema.fields ++
Array(
StructField("ModelValidityPeriod", StringType, false),
StructField("ModelValidityPeriod2", StringType, false)
)
)
spark.sqlContext.createDataFrame(raw_customer_df.rdd.mapPartitions(transformRows), newSchema).show
【问题讨论】:
-
为什么不使用 spark-redis 连接器? (github.com/RedisLabs/spark-redis)
-
嗨,Korland,好的,我们可以,但这不是问题。主要关注如何在 mappartitions 中使用 Redis 进行查找时访问数据帧行和列。非常感谢任何帮助。
标签: apache-spark redis streaming lookup