【问题标题】:Apache Spark nested iterations in Scala to generate stats RDDApache Spark 在 Scala 中嵌套迭代以生成统计数据 RDD
【发布时间】:2016-10-10 20:58:29
【问题描述】:

我有一个由 rowkey = client_id、campaign = {campaign_id:campaign_name} 的 Json 数组制作的 RDD

val clientsRDD = resultRDD.map(ClientRow.parseClientRow)
// change  RDD of ClientRow  objects to a DataFrame
val clientsDF = clientsRDD.toDF()
// Return the schema of this DataFrame
clientsDF.printSchema()
// print each line DataFrame
clientsDF.collect().foreach(println)

输出:

root
 |-- rowkey: string (nullable = true)
 |-- campaigns: string (nullable = true)

[1,[{"1000":"campaign1"},{"1001":"campaign2"}]]
[2,[{"1002":"campaign3"}]]

我还有一个 RDD,其中包含来自 HBase 的所有客户和活动数据的记录。

记录RDD

rowkey                 type         body
client_id-campaign_id, record_type, record_text 

我的目标是为每个客户(考虑其所有活动)和每个活动生成统计信息,例如计算所有 client_id 记录、按类型分组和计算每个活动记录、按类型分组。

client1
records:100, login:20, actions:80

client1 campaign1  
records:70, login:16, actions:50

client1 campaign2
records:30, login:4, actions:30

最后我想写统计数据。

使用 Scala 在 Spark 中执行此操作的最佳方法是什么? 我是否必须迭代clientsRDD(映射?),并为每一行生成不同的RDD映射recordsRDD?

【问题讨论】:

    标签: scala apache-spark mapreduce apache-spark-sql rdd


    【解决方案1】:

    首先,您需要为活动字段定义架构: 它的意思是 您使用

    定义架构
    val schema = StructType(Seq(StructField("rowkey", StringType, true),
    StructField("campaigns", StructType(
      StructField("id", StringType, true) ::
        StructField("name", StringType, true) :: Nil
    ))
    

    ))

    然后您可以在活动字段中使用explode 方法使行变平。

    val df = sqlContext.createDataFrame(clientsRDD, schema)
    df.select(col("rowkey"), explode(col("campaigns")).as("campaign")).filter(col("campaign.id") === 1)
    

    【讨论】:

      猜你喜欢
      • 2015-06-14
      • 2014-07-12
      • 1970-01-01
      • 1970-01-01
      • 2019-06-25
      • 1970-01-01
      • 1970-01-01
      • 2015-10-26
      • 1970-01-01
      相关资源
      最近更新 更多