【问题标题】:Data transfer from csv format to Redis hash format in DatabricksDatabricks 中从 csv 格式到 Redis 哈希格式的数据传输
【发布时间】:2020-11-09 21:10:30
【问题描述】:

我有一个 Azure 系统,分为三个部分:

  1. 我有一些 csv 文件的 Azure 数据湖存储。
  2. Azure Databricks 我需要进行一些处理 - 正是将该 csv 文件转换为 Redis 哈希格式。
  3. Azure Redis 缓存我应该将转换后的数据放在哪里。

在 databricks 文件系统中挂载存储后,需要处理一些数据。 如何将位于 databricks 文件系统中的 csv 数据转换为 redisHash 格式并正确地放入 Redis? 具体来说,我不确定如何通过下面的代码进行正确的映射。或者也许有一些我找不到的额外转移到 SQL 表的方法。

这是我在 scala 上编写的代码示例:

import com.redislabs.provider.redis._

val redisServerDnsAddress = "HOST"
val redisPortNumber = 6379
val redisPassword = "Password"
val redisConfig = new RedisConfig(new RedisEndpoint(redisServerDnsAddress, redisPortNumber, redisPassword))


val data = spark.read.format("com.databricks.spark.csv").option("header", "true").option("inferSchema", "true").load("/mnt/staging/data/file.csv")

// What is the right way of mapping?
val ds = table("data").select("Prop1", "Prop2", "Prop3", "Prop4", "Prop5" ).distinct.na.drop().map{x =>
  (x.getString(0), x.getString(1), x.getString(2), x.getString(3), x.getString(4))
}

sc.toRedisHASH(ds, "data")

错误:

error: type mismatch;
 found   : org.apache.spark.sql.Dataset[(String, String)]
 required: org.apache.spark.rdd.RDD[(String, String)]
sc.toRedisHASH(ds, "data")

如果我这样写最后一串代码:

sc.toRedisHASH(ds.rdd, "data")

错误:

org.apache.spark.sql.AnalysisException: Table or view not found: data;

【问题讨论】:

  • 当您尝试查询表或视图时会发生该错误。您已经从 csv 构建了数据框,根据 REDIS 连接器文档将其转换为 RDD。

标签: scala apache-spark redis databricks azure-databricks


【解决方案1】:

准备一些示例数据来模拟从 CSV 文件加载的数据。

    val rdd = spark.sparkContext.parallelize(Seq(Row("1", "2", "3", "4", "5", "6", "7")))
    val structType = StructType(
      Seq(
        StructField("Prop1", StringType),
        StructField("Prop2", StringType),
        StructField("Prop3", StringType),
        StructField("Prop4", StringType),
        StructField("Prop5", StringType),
        StructField("Prop6", StringType),
        StructField("Prop7", StringType)
      )
    )
    val data = spark.createDataFrame(rdd, structType)

转换:

val transformedData = data.select("Prop1", "Prop2", "Prop3", "Prop4", "Prop5").distinct.na.drop()

将数据框写入 Redis,使用Prop1 作为键,data 作为 Redis 表名。见docs

    transformedData
      .write
      .format("org.apache.spark.sql.redis")
      .option("key.column", "Prop1")
      .option("table", "data")
      .mode(SaveMode.Overwrite)
      .save()

查看 Redis 中的数据:

127.0.0.1:6379> keys data:*
1) "data:1"

127.0.0.1:6379> hgetall data:1
1) "Prop5"
2) "5"
3) "Prop2"
4) "2"
5) "Prop4"
6) "4"
7) "Prop3"
8) "3"

【讨论】:

  • 非常感谢您的回答。似乎一切都应该工作,但我无法在文档中找到通过您的代码示例将我的配置连接到 redis 的位置。据我了解,如果没有此配置,则会发生以下错误:redis.clients.jedis.exceptions.JedisConnectionException: Could not get a resource from the pool
  • 您应该将它们作为 spark 配置选项,例如val spark = SparkSession .builder() .appName("redis-df") .master("local[*]") .config("spark.redis.host", "localhost") .config("spark.redis. port", "6379") .getOrCreate() 另一种方法是直接使用数据框选项覆盖这些连接设置,请参阅github.com/RedisLabs/spark-redis/blob/master/doc/…
  • 在 Databricks 中,可以更改 spark 配置选项,如下所述stackoverflow.com/questions/58688544/…
  • 是的,使用覆盖数据框的选项效果很好。再次感谢!
  • @Maksim 很高兴知道您的问题已解决。您可以接受它作为答案(单击答案旁边的复选标记以将其从灰色切换为已填充)。这对其他社区成员可能是有益的。谢谢。
猜你喜欢
  • 1970-01-01
  • 2018-07-04
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-06-05
  • 2015-09-18
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多