2017-07-19 84 views
1

我有要求將流數據加載到DynamoDB表中的要求。我嘗試了下面的代碼。無法在spark中創建DynamoDB客戶端執行程序

object UnResolvedLoad { 

    def main(args: Array[String]){ 
    val spark = SparkSession.builder().appName("unresolvedload").enableHiveSupport().getOrCreate() 
    val tokensDf = spark.sql("select * from unresolved_logic.unresolved_dynamo_load") 
    tokensDf.foreachPartition { x => loadFunc(x) } 
    } 


    def loadFunc(iter : Iterator[org.apache.spark.sql.Row]) = { 

     val client:AmazonDynamoDB = AmazonDynamoDBClientBuilder.standard().build() 
     val dynamoDB:DynamoDB = new DynamoDB(client) 
     val table:Table = dynamoDB.getTable("UnResolvedTokens") 

     while(iter.hasNext){ 
     val cur = iter.next() 
     val item:Item = new Item().withString("receiverId ", cur.get(2).asInstanceOf[String]). 
       withString("payload_id", cur.get(0).asInstanceOf[String]). 
       withString("payload_confirmation_code", cur.get(1).asInstanceOf[String]). 
       withString("token", cur.get(3).asInstanceOf[String]) 

     table.putItem(item) 

     } 

} 

}

當我執行火花提交它不能夠實例類。以下是錯誤信息。它說它不能實例化Class。幫助表示讚賞。 有沒有一種方法,我們可以節省星火DataSet轉換亞馬遜DynamoDB

, executor 5): java.lang.NoClassDefFoundError: Could not initialize class com.amazonaws.services.dynamodbv2.AmazonDynamoDBClientBuilder 
     at com.dish.payloads.UnResolvedLoad$.loadFunc(UnResolvedLoad.scala:22) 
     at com.dish.payloads.UnResolvedLoad$$anonfun$main$1.apply(UnResolvedLoad.scala:16) 
     at com.dish.payloads.UnResolvedLoad$$anonfun$main$1.apply(UnResolvedLoad.scala:16) 
     at org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1$$anonfun$apply$29.apply(RDD.scala:926) 
     at org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1$$anonfun$apply$29.apply(RDD.scala:926) 
     at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:1951) 
     at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:1951) 
     at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87) 
     at org.apache.spark.scheduler.Task.run(Task.scala:99) 
     at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:322) 
     at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142) 
     at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617) 
     at java.lang.Thread.run(Thread.java:748) 

17/07/19 17:35:15 INFO TaskSetManager: Lost task 26.0 in stage 0.0 (TID 26) on ip-10-176-225-151.us-west-2.compute.internal, executor 5: java.lang.NoClassDefFoundError (Could not initialize class com.amazonaws.services.dynamodbv2.AmazonDynamoDBClientBuilder) [duplicate 1] 
17/07/19 17:35:15 WARN TaskSetManager: Lost task 6.0 in stage 0.0 (TID 6, ip-10-176-225-151.us-west-2.compute.internal, executor 5): java.lang.IllegalAccessError: tried to access class com.amazonaws.services.dynamodbv2.AmazonDynamoDBClientConfigurationFactory from class com.amazonaws.services.dynamodbv2.AmazonDynamoDBClientBuilder 
     at com.amazonaws.services.dynamodbv2.AmazonDynamoDBClientBuilder.<clinit>(AmazonDynamoDBClientBuilder.java:30) 
     at com.dish.payloads.UnResolvedLoad$.loadFunc(UnResolvedLoad.scala:22) 
     at com.dish.payloads.UnResolvedLoad$$anonfun$main$1.apply(UnResolvedLoad.scala:16) 
     at com.dish.payloads.UnResolvedLoad$$anonfun$main$1.apply(UnResolvedLoad.scala:16) 
     at org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1$$anonfun$apply$29.apply(RDD.scala:926) 
     at org.apache.spark.rdd.RDD$$anonfun$foreachPartition$1$$anonfun$apply$29.apply(RDD.scala:926) 
     at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:1951) 
     at org.apache.spark.SparkContext$$anonfun$runJob$5.apply(SparkContext.scala:1951) 
     at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87) 
     at org.apache.spark.scheduler.Task.run(Task.scala:99) 
     at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:322) 
     at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142) 
     at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617) 
     at java.lang.Thread.run(Thread.java:748) 

回答

1

我終於能夠通過使用DynamoDB API的較低版本來解決它。 EMR 5.7僅支持1.10.75.1。以下是對我來說工作正常的代碼。

object UnResolvedLoad { 

    def main(args: Array[String]){ 
    val spark = SparkSession.builder().appName("unresolvedload").enableHiveSupport().getOrCreate() 
    val tokensDf = spark.sql("select * from unresolved_logic.unresolved_dynamo_load") 
    tokensDf.foreachPartition { x => loadFunc(x) } 
    } 


    def loadFunc(iter : Iterator[org.apache.spark.sql.Row]) = { 

     val client:AmazonDynamoDBClient = new AmazonDynamoDBClient(); 
     val usWest2 = Region.getRegion(Regions.US_WEST_2); 
     client.setRegion(usWest2) 



     while(iter.hasNext){ 
     val cur = iter.next() 

     val putMap = Map("receiverId" -> new AttributeValue(cur.get(2).asInstanceOf[String]), 
          "payload_id" -> new AttributeValue(cur.get(0).asInstanceOf[String]), 
          "payload_confirmation_code" -> new AttributeValue(cur.get(1).asInstanceOf[String]), 
          "token" -> new AttributeValue(cur.get(3).asInstanceOf[String])).asJava 

     val putItemRequest:PutItemRequest = new PutItemRequest("UnResolvedTokens",putMap) 
     client.putItem(putItemRequest) 
     } 

    } 
} 
+0

爲我節省了大量時間!謝啦。你是怎麼想到的?真的很難想到這個方向 – KAs

+0

我從AWS支持論壇得到了答案 –

相關問題