【问题标题】:Spark 2.2.0 - How to write/read DataFrame to DynamoDBSpark 2.2.0 - 如何将 DataFrame 写入/读取到 DynamoDB
【发布时间】:2018-05-23 04:31:53
【问题描述】:

我希望我的 Spark 应用程序从 DynamoDB 读取表,执行操作,然后将结果写入 DynamoDB。

将表格读入 DataFrame

现在,我可以将 DynamoDB 中的表作为 hadoopRDD 读取到 Spark 中,并将其转换为 DataFrame。但是,我不得不使用正则表达式从AttributeValue 中提取值。有没有更好/更优雅的方式?在 AWS API 中找不到任何内容。

package main.scala.util

import org.apache.spark.sql.SparkSession
import org.apache.spark.SparkContext
import org.apache.spark.sql.SQLContext
import org.apache.spark.sql.functions._
import org.apache.spark.sql.types._
import org.apache.spark.rdd.RDD
import scala.util.matching.Regex
import java.util.HashMap

import com.amazonaws.services.dynamodbv2.model.AttributeValue
import org.apache.hadoop.io.Text;
import org.apache.hadoop.dynamodb.DynamoDBItemWritable
/* Importing DynamoDBInputFormat and DynamoDBOutputFormat */
import org.apache.hadoop.dynamodb.read.DynamoDBInputFormat
import org.apache.hadoop.dynamodb.write.DynamoDBOutputFormat
import org.apache.hadoop.mapred.JobConf
import org.apache.hadoop.io.LongWritable

object Tester {

  // {S: 298905396168806365,} 
  def extractValue : (String => String) = (aws:String) => {
    val pat_value = "\\s(.*),".r

    val matcher = pat_value.findFirstMatchIn(aws)
                matcher match {
                case Some(number) => number.group(1).toString
                case None => ""
        }
  }


   def main(args: Array[String]) {
    val spark = SparkSession.builder().getOrCreate()
    val sparkContext = spark.sparkContext

      import spark.implicits._

      // UDF to extract Value from AttributeValue 
      val col_extractValue = udf(extractValue)

  // Configure connection to DynamoDB
  var jobConf_add = new JobConf(sparkContext.hadoopConfiguration)
      jobConf_add.set("dynamodb.input.tableName", "MyTable")
      jobConf_add.set("dynamodb.output.tableName", "MyTable")
      jobConf_add.set("mapred.output.format.class", "org.apache.hadoop.dynamodb.write.DynamoDBOutputFormat")
      jobConf_add.set("mapred.input.format.class", "org.apache.hadoop.dynamodb.read.DynamoDBInputFormat")


      // org.apache.spark.rdd.RDD[(org.apache.hadoop.io.Text, org.apache.hadoop.dynamodb.DynamoDBItemWritable)]
      var hadooprdd_add = sparkContext.hadoopRDD(jobConf_add, classOf[DynamoDBInputFormat], classOf[Text], classOf[DynamoDBItemWritable])

      // Convert HadoopRDD to RDD
      val rdd_add: RDD[(String, String)] = hadooprdd_add.map {
      case (text, dbwritable) => (dbwritable.getItem().get("PIN").toString(), dbwritable.getItem().get("Address").toString())
      }

      // Convert RDD to DataFrame and extract Values from AttributeValue
      val df_add = rdd_add.toDF()
                  .withColumn("PIN", col_extractValue($"_1"))
                  .withColumn("Address", col_extractValue($"_2"))
                  .select("PIN","Address")
   }
}

将 DataFrame 写入 DynamoDB

stackoverflow 和其他地方的许多答案只指向blog post 和emr-dynamodb-hadoop github。这些资源都没有真正演示如何写入 DynamoDB。

I tried converting 我的DataFrame 到RDD[Row] 不成功。

df_add.rdd.saveAsHadoopDataset(jobConf_add)

将此 DataFrame 写入 DynamoDB 的步骤是什么? (如果你告诉我如何控制overwrite 和putItem ,可以获得奖励积分;)

注意:df_add 与 DynamoDB 中的 MyTable 具有相同的架构。

编辑:我遵循this answer 的建议,该建议指向Using Spark SQL for ETL 上的这篇文章:

// Format table to DynamoDB format
  val output_rdd =  df_add.as[(String,String)].rdd.map(a => {
    var ddbMap = new HashMap[String, AttributeValue]()

    // Field PIN
    var PINValue = new AttributeValue() // New AttributeValue
    PINValue.setS(a._1)                 // Set value of Attribute as String. First element of tuple
    ddbMap.put("PIN", PINValue)         // Add to HashMap

    // Field Address
    var AddValue = new AttributeValue() // New AttributeValue
    AddValue.setS(a._2)                 // Set value of Attribute as String
    ddbMap.put("Address", AddValue)     // Add to HashMap

    var item = new DynamoDBItemWritable()
    item.setItem(ddbMap)

    (new Text(""), item)
  })             

  output_rdd.saveAsHadoopDataset(jobConf_add) 

但是,尽管遵循了文档,但现在我得到了java.lang.ClassCastException: java.lang.String cannot be cast to org.apache.hadoop.io.Text ...您有什么建议吗?

编辑 2:在Using Spark SQL for ETL 上仔细阅读这篇文章:

在您拥有 DataFrame 后,执行转换以获得与 DynamoDB 自定义输出格式知道如何编写的类型相匹配的 RDD。自定义输出格式需要一个包含 Text 和 DynamoDBItemWritable 类型的元组。

考虑到这一点,下面的代码正是 AWS 博客文章所建议的,除了我将 output_df 转换为 rdd 否则 saveAsHadoopDataset 不起作用。现在,我收到了Exception in thread "main" scala.reflect.internal.Symbols$CyclicReference: illegal cyclic reference involving object InterfaceAudience。我已经走到了尽头!

      // Format table to DynamoDB format
  val output_df =  df_add.map(a => {
    var ddbMap = new HashMap[String, AttributeValue]()

    // Field PIN
    var PINValue = new AttributeValue() // New AttributeValue
    PINValue.setS(a.get(0).toString())                 // Set value of Attribute as String
    ddbMap.put("PIN", PINValue)         // Add to HashMap

    // Field Address
    var AddValue = new AttributeValue() // New AttributeValue
    AddValue.setS(a.get(1).toString())                 // Set value of Attribute as String
    ddbMap.put("Address", AddValue)     // Add to HashMap

    var item = new DynamoDBItemWritable()
    item.setItem(ddbMap)

    (new Text(""), item)
  })             

  output_df.rdd.saveAsHadoopDataset(jobConf_add)   

【问题讨论】:

  • 我遇到了类似的 scala.reflect.internal.Symbols$CyclicReference 错误,请帮忙?
  • 我没有找到解决方案,而是使用了 S3。
  • 这是从 DynamoDB 读取的更简单的方法,val simple2: RDD[(String)] = data.map { case (text, dbwritable) => (dbwritable.toString)} 然后触发。 read.json(simple2).registerTempTable("gooddata") 然后 spark.sql("select replace(replace(split(cast(address as string),',')[0],']',''), '[','') as housenumber from gooddata").show(false)
  • @Béatrice Moissinac 你现在有解决方案吗...如果可能,请分享 github 链接

标签: scala apache-spark amazon-dynamodb amazon-emr


【解决方案1】:

我正在关注“使用 Spark SQL 进行 ETL”链接,发现同样的“非法循环引用”异常。 该异常的解决方案非常简单(但我花了 2 天时间才弄清楚),如下所示。关键是在dataframe的RDD上使用map函数,而不是dataframe本身。

val ddbConf = new JobConf(spark.sparkContext.hadoopConfiguration)
ddbConf.set("dynamodb.output.tableName", "<myTableName>")
ddbConf.set("dynamodb.throughput.write.percent", "1.5")
ddbConf.set("mapred.input.format.class", "org.apache.hadoop.dynamodb.read.DynamoDBInputFormat")
ddbConf.set("mapred.output.format.class", "org.apache.hadoop.dynamodb.write.DynamoDBOutputFormat")


val df_ddb =  spark.read.option("header","true").parquet("<myInputFile>")
val schema_ddb = df_ddb.dtypes

var ddbInsertFormattedRDD = df_ddb.rdd.map(a => {
    val ddbMap = new HashMap[String, AttributeValue]()

    for (i <- 0 to schema_ddb.length - 1) {
        val value = a.get(i)
        if (value != null) {
            val att = new AttributeValue()
            att.setS(value.toString)
            ddbMap.put(schema_ddb(i)._1, att)
        }
    }

    val item = new DynamoDBItemWritable()
    item.setItem(ddbMap)

    (new Text(""), item)
}
)

ddbInsertFormattedRDD.saveAsHadoopDataset(ddbConf)

【讨论】:

【解决方案2】:

这是一个更简单的工作示例。

例如使用 Hadoop RDD 从 Kinesis Stream 写入 DynamoDB:-

https://github.com/kali786516/Spark2StructuredStreaming/blob/master/src/main/scala/com/dataframe/part11/kinesis/consumer/KinesisSaveAsHadoopDataSet/TransactionConsumerDstreamToDynamoDBHadoopDataSet.scala

用于使用 Hadoop RDD 从 DynamoDB 读取并使用 spark SQL 而不使用正则表达式。

val ddbConf = new JobConf(spark.sparkContext.hadoopConfiguration)
    //ddbConf.set("dynamodb.output.tableName", "student")
    ddbConf.set("dynamodb.input.tableName", "student")
    ddbConf.set("dynamodb.throughput.write.percent", "1.5")
    ddbConf.set("dynamodb.endpoint", "dynamodb.us-east-1.amazonaws.com")
    ddbConf.set("dynamodb.regionid", "us-east-1")
    ddbConf.set("dynamodb.servicename", "dynamodb")
    ddbConf.set("dynamodb.throughput.read", "1")
    ddbConf.set("dynamodb.throughput.read.percent", "1")
    ddbConf.set("mapred.input.format.class", "org.apache.hadoop.dynamodb.read.DynamoDBInputFormat")
    ddbConf.set("mapred.output.format.class", "org.apache.hadoop.dynamodb.write.DynamoDBOutputFormat")
    //ddbConf.set("dynamodb.awsAccessKeyId", credentials.getAWSAccessKeyId)
    //ddbConf.set("dynamodb.awsSecretAccessKey", credentials.getAWSSecretKey)


val data = spark.sparkContext.hadoopRDD(ddbConf, classOf[DynamoDBInputFormat], classOf[Text], classOf[DynamoDBItemWritable])

val simple2: RDD[(String)] = data.map { case (text, dbwritable) => (dbwritable.toString)}

spark.read.json(simple2).registerTempTable("gooddata")

spark.sql("select replace(replace(split(cast(address as string),',')[0],']',''),'[','') as housenumber from gooddata").show(false)

【讨论】:

  • 您的链接是 404
猜你喜欢
  • 2017-04-03
  • 2020-06-12
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-04-20
  • 1970-01-01
  • 2019-09-20
  • 2018-03-09
相关资源
最近更新 更多