【发布时间】:2018-09-01 09:02:27
【问题描述】:
我有一个具有以下架构的 spark 数据帧,并尝试使用 Avro 将此数据帧流式传输到 Kafka
```root
|-- clientTag: struct (nullable = true)
| |-- key: string (nullable = true)
|-- contactPoint: struct (nullable = true)
| |-- email: string (nullable = true)
| |-- type: string (nullable = true)
|-- performCheck: string (nullable = true)```
样本记录:{"performCheck" : "N", "clientTag" :{"key":"value"}, "contactPoint": {"email":"abc@gmail.com", "type":"EML"}}
Avro 架构:
{
"name":"Message",
"namespace":"kafka.sample.avro",
"type":"record",
"fields":[
{"type":"string", "name":"id"},
{"type":"string", "name":"email"}
{"type":"string", "name":"type"}
]
}
我有几个问题。
- 将
org.apache.spark.sql.Row转换为Avro 消息的最佳方法是什么,因为我想从每一行的数据框中提取email和type并使用这些值来构造Avro 消息? - 最终,所有的 Avro 消息都将发送到 Kafka。那么,如果生产过程中出现错误,如何将所有未能生产的 Row 收集到 Kafka 并返回数据帧?
感谢您的帮助
【问题讨论】:
标签: apache-spark dataframe apache-kafka spark-dataframe spark-streaming