【发布时间】: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() 创建合适的数据框
【问题讨论】:
-
您是否尝试进行搜索? spark-streaming + dataframe
-
它没有帮助,因为我是 scala 的新手。我无法弄清楚如何将 avro[String,String] 转换为数据帧stackoverflow.com/questions/41237929/…
-
这是我的答案谢谢 Maasg 找到了答案
标签: dataframe apache-kafka spark-streaming