【问题标题】:Spark Streaming reach dataframe columns and add new column looking up to RedisSpark Streaming 到达数据框列并添加查找 Redis 的新列
【发布时间】: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


【解决方案1】:

迭代器iter 表示数据帧行上的迭代器。因此,如果我正确地回答了您的问题,您可以通过迭代 iter 并调用来访问列值

row.getAs[Column_Type](column_name)

类似的东西

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
    val res = iter.map { row =>
      val columnValue = row.getAs[String]("column_name")
      // lookup in redis
      val valueFromRedis = redisConn.get(...)
      Row.fromSeq(row.toSeq ++ Array[Any](valueFromRedis))
    }.toList

    redisConn.close()
    res.iterator
  }

【讨论】:

  • 嗨 fe2s,感谢您的建议。这个动态添加字段的架构可以做什么?
  • 另外,在原始数据框中有嵌套列。 val columnValue = row.getAs[String]("C.field1") 这在到达嵌套列时出错
  • 嗨@mustangc。不明白问题重新分级架构。至于访问嵌套列,应该可以使用类似 row.getAs[Row]("C").getAs[String]("field1")
  • 嗨,@fe2s。我修复了嵌套项目解析的问题,如下所示: var BirthPlaceCityCode = Try( row.getAs[Row]("ObjectLevel1") .getAs[GenericRowWithSchema]("ObjectLevel2") .getAs[Integer]("BirthPlaceCityCode")).getOrElse (-1).toString 这没关系。我的主要目的是:a)来自kafka主题的ReadStream b)丰富withRedis并“添加新列” c)WriteStream到kafkatopic“带有新列”所以,在执行WriteStream时我应该使用foreachbatch吗?特别感谢
猜你喜欢
  • 2015-11-19
  • 1970-01-01
  • 2021-12-07
  • 1970-01-01
  • 2017-06-16
  • 2017-05-20
  • 2018-05-10
  • 2011-05-12
  • 2017-05-18
相关资源
最近更新 更多