【问题标题】:inappropriate output while creating a dataframe创建数据框时输出不适当
【发布时间】:2016-12-21 00:50:46
【问题描述】:

我正在尝试使用 scala 应用程序从 kafka 主题流式传输数据。我能够从主题中获取数据,但是如何从中创建数据框?

这是数据(字符串,字符串格式)

{
  "action": "AppEvent",
  "tenantid": 298,
  "lat": 0.0,
  "lon": 0.0,
  "memberid": 16390,
  "event_name": "CATEGORY_CLICK",
  "productUpccd": 0,
  "device_type": "iPhone",
  "device_os_ver": "10.1",
  "item_name": "CHICKEN"
}

我尝试了几种方法来做到这一点,但都没有产生令人满意的结果。

 +--------------------+ |                  _1|
 +--------------------+ |{"action":"AppEve...| |{"action":"AppEve...| |{"action":"AppEve...| |{"action":"AppEve...| |{"action":"AppEve...|
 |{"action":"AppEve...| |{"action":"AppEve...| |{"action":"AppEve...|
 |{"action":"AppEve...| |{"action":"AppEve...|

谁能告诉如何进行映射,以便每个字段像表格一样进入单独的列。数据为avro格式。

这是从主题中获取数据的代码。

val ssc = new StreamingContext(sc, Seconds(2))
val kafkaConf = Map[String, String]("metadata.broker.list" -> "####",
     "zookeeper.connect" -> "########",
     "group.id" -> "KafkaConsumer",
     "zookeeper.connection.timeout.ms" -> "1000000")
val topicMaps = Map("fishbowl" -> 1)
val messages  = KafkaUtils.createStream[String, String,DefaultDecoder, DefaultDecoder](ssc, kafkaConf, topicMaps, StorageLevel.MEMORY_ONLY_SER).map(_._2)

请指导我如何使用 foreachRDD func 和 map() 创建合适的数据框

【问题讨论】:

标签: dataframe apache-kafka spark-streaming


【解决方案1】:

从 rdd 中创建数据框,而不管其案例类模式如何。 使用下面的逻辑

stream.foreachRDD(
  rdd => {
     val dataFrame = sqlContext.read.json(rdd.map(_._2)) 
dataFrame.show()
        })

这里的流是从 kafkaUtils.createStream() 创建的 rdd

【讨论】:

  • 干得好。关于“无论其格式或案例类模式如何”的评论,这并不完全正确 => 这仅适用于 JSON 格式的记录。
  • @maasg 谢谢先生,编辑了我的评论。当我用 avro 解决它时(它的模式仍然在 json 中)
猜你喜欢
  • 1970-01-01
  • 2015-08-29
  • 2020-12-30
  • 2021-05-15
  • 1970-01-01
  • 1970-01-01
  • 2021-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多