【问题标题】:Spark redis connector to write data into specific index of the redisSpark redis连接器将数据写入redis的特定索引
【发布时间】:2020-07-08 06:10:03
【问题描述】:

我正在尝试从 Cassandra 读取数据并写入特定索引的 Redis。假设是 Redis DB 5。

我需要将所有数据以 hashmap 格式写入 Redis DB 索引 5。

 val spark = SparkSession.builder()
  .appName("redis-df")
  .master("local[*]")
  .config("spark.redis.host", "localhost")
  .config("spark.redis.port", "6379")
  .config("spark.redis.db", 5)
  .config("spark.cassandra.connection.host", "localhost")
  .getOrCreate()

  import spark.implicits._
    val someDF = Seq(
      (8, "bat"),
      (64, "mouse"),
      (-27, "horse")
    ).toDF("number", "word")

    someDF.write
      .format("org.apache.spark.sql.redis")
      .option("keys.pattern", "*")
      //.option("table", "person"). // Is it mandatory ?
      .save()

我可以在没有表名的情况下将数据保存到 Redis 中吗?实际上只是我想将所有数据保存到没有表名的 Redis 索引 5 中是否可能? 我已经浏览了 spark Redis 连接器的文档,我没有看到任何与此相关的示例。 文档链接:https://github.com/RedisLabs/spark-redis/blob/master/doc/dataframe.md#writing

我目前正在使用这个版本的 spark redis-connector

    <dependency>
        <groupId>com.redislabs</groupId>
        <artifactId>spark-redis_2.11</artifactId>
        <version>2.5.0</version>
    </dependency>

有人遇到过这个问题吗?有什么解决办法吗?

如果我没有在配置中提及表名,我得到的错误

失败

  java.lang.IllegalArgumentException: Option 'table' is not set.
  at org.apache.spark.sql.redis.RedisSourceRelation$$anonfun$tableName$1.apply(RedisSourceRelation.scala:208)
  at org.apache.spark.sql.redis.RedisSourceRelation$$anonfun$tableName$1.apply(RedisSourceRelation.scala:208)
  at scala.Option.getOrElse(Option.scala:121)
  at org.apache.spark.sql.redis.RedisSourceRelation.tableName(RedisSourceRelation.scala:208)
  at org.apache.spark.sql.redis.RedisSourceRelation.saveSchema(RedisSourceRelation.scala:245)
  at org.apache.spark.sql.redis.RedisSourceRelation.insert(RedisSourceRelation.scala:121)
  at org.apache.spark.sql.redis.DefaultSource.createRelation(DefaultSource.scala:30)
  at org.apache.spark.sql.execution.datasources.SaveIntoDataSourceCommand.run(SaveIntoDataSourceCommand.scala:45)
  at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:70)
  at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:68)

【问题讨论】:

    标签: scala dataframe apache-spark pyspark redis


    【解决方案1】:

    table 选项是强制性的。这个想法是您指定表名,因此可以从提供该表名的 Redis 读取数据帧。 在您的情况下,另一种选择是将数据帧转换为键/值 RDD 并使用 sc.toRedisKV(rdd)

    【讨论】:

      【解决方案2】:

      我不得不不同意。我正在处理与您完全相同的问题。这是我发现的:

      1. 您必须引用表或键模式。 (例如)

        df = spark.read.format("org.apache.spark.sql.redis")
        .option("keys.pattern", "rec-*")
        .option("infer.schema", True).load()

      在我的例子中,我使用的是 HASH,并且 HASH 键都以“rec-”开头,后跟一个 int。 spark-redis 代码将“rec-”视为一个表。如前所述,诀窍是如果您想将数据读回 Spark。它需要一个表名,但似乎使用冒号作为分隔符。由于我想做读/写,我只是将我的表名更改为“rec:”并且很好。

      我认为您的困惑源于在您的示例中,您在 Spark 中只定义了一条记录。如果你有两个呢? Redis 需要创建两个不同的键,例如“person:1”或“person:2”。它使用术语表来描述“人”。是钥匙还是桌子?文档似乎不一致。

      我目前的问题是能够通过以某种方式更改数据库上下文 .config("spark.redis.db", 5) 来保存到不同的 Redis 数据库。当我在 df.write.format 中使用它时,这似乎对我不起作用。有什么想法吗?

      【讨论】:

        猜你喜欢
        • 2019-06-17
        • 1970-01-01
        • 1970-01-01
        • 2020-06-07
        • 2013-06-04
        • 1970-01-01
        • 2020-03-08
        • 2021-11-11
        • 2022-11-07
        相关资源
        最近更新 更多