有没有使用模式转换方式的Avro从消息卡夫卡与火花到数据帧?用户记录的模式文件:
{
"fields": [
{ "name": "firstName", "type": "string" },
{ "name": "lastName", "type": "string" }
],
"name": "user",
"type": "record"
}
Run Code Online (Sandbox Code Playgroud)
来自SqlNetworkWordCount示例和Kafka,Spark和Avro的代码片段- 第3部分,生成和使用Avro消息来读取消息.
object Injection {
val parser = new Schema.Parser()
val schema = parser.parse(getClass.getResourceAsStream("/user_schema.json"))
val injection: Injection[GenericRecord, Array[Byte]] = GenericAvroCodecs.toBinary(schema)
}
...
messages.foreachRDD((rdd: RDD[(String, Array[Byte])]) => {
val sqlContext = SQLContextSingleton.getInstance(rdd.sparkContext)
import sqlContext.implicits._
val df = rdd.map(message => Injection.injection.invert(message._2).get)
.map(record => User(record.get("firstName").toString, records.get("lastName").toString)).toDF()
df.show() …Run Code Online (Sandbox Code Playgroud) 在Row Java API中有一个row.schema(),但是没有row.set(StructType模式).
我也试过RowFactorie.create(objets),但我不知道如何继续
更新:
问题是当我修改示例中的工作者的结构时如何生成新的数据帧
DataFrame sentenceData = jsql.createDataFrame(jrdd, schema);
List<Row> resultRows2 = sentenceData.toJavaRDD()
.map(new MyFunction<Row, Row>(parameters) {
/** my map function **//
public Row call(Row row) {
// I want to change Row definition adding new columns
Row newRow = functionAddnewNewColumns (row);
StructType newSchema = functionGetNewSchema (row.schema);
// Here I want to insert the structure
//
return newRow
}
}
}).collect();
JavaRDD<Row> jrdd = jsc.parallelize(resultRows);
// Here is the problema I don't know how to get the new …Run Code Online (Sandbox Code Playgroud) 这实际上与我之前的问题相同,但使用Avro而不是JSON作为数据格式.
我正在使用Spark数据帧,它可以从几个不同的模式版本之一加载数据:
// Version One
{"namespace": "com.example.avro",
"type": "record",
"name": "MeObject",
"fields": [
{"name": "A", "type": ["null", "int"], "default": null}
]
}
// Version Two
{"namespace": "com.example.avro",
"type": "record",
"name": "MeObject",
"fields": [
{"name": "A", "type": ["null", "int"], "default": null},
{"name": "B", "type": ["null", "int"], "default": null}
]
}
Run Code Online (Sandbox Code Playgroud)
我正在使用Spark Avro加载数据.
DataFrame df = context.read()
.format("com.databricks.spark.avro")
.load("path/to/avro/file");
Run Code Online (Sandbox Code Playgroud)
可以是Version One文件或Version Two文件.但是我希望能够以相同的方式处理它,将未知值设置为"null".我之前的问题中的建议是设置模式,但是我不想重复自己在.avro文件和火花StructType和朋友中编写模式.如何将avro架构(文本文件或生成的MeObject.getClassSchema())转换为火花StructType?
Spark Avro有一个SchemaConverters,但它都是私有的,并返回一些奇怪的内部对象.
我从Avaf(序列化器和反序列化器)获得了kafka主题的推文.然后我创建了一个Spark消费者,它在RDD [GenericRecord]的Dstream中提取推文.现在我想将每个rdd转换为数据帧,以通过SQL分析这些推文.有什么解决方案可以将RDD [GenericRecord]转换为数据帧吗?