您可以尝试这种方式(Spark 1.6)。
people.csv
Michael, 29
Andy, 30
Justin, 19
派斯帕克
file = sc.textFile("people.csv")
df = file.map(lambda line: line.split(',')).toDF(['name','age'])
>>> df.show()
+-------+---+
| name|age|
+-------+---+
|Michael| 29|
| Andy| 30|
| Justin| 19|
+-------+---+
df.write.format("com.databricks.spark.avro").save("peopleavro")
人脉
{u'age': u' 29', u'name': u'Michael'}
{u'age': u' 30', u'name': u'Andy'}
{u'age': u' 19', u'name': u'Justin'}
如果你需要维护数据类型,然后创建一个模式并传递它。
schema = StructType([StructField("name",StringType(),True),StructField("age",IntegerType(),True)])
df = file.map(lambda line: line.split(',')).toDF(schema)
>>> df.printSchema()
root
|-- name: string (nullable = true)
|-- age: integer (nullable = true)
现在你的 avro 有
{
"type" : "record",
"name" : "topLevelRecord",
"fields" : [ {
"name" : "name",
"type" : [ "string", "null" ]
}, {
"name" : "age",
"type" : [ "int", "null" ]
} ]
}