我正在尝试从SparkSQL中写入和读取镶木面板文件。出于模式演变的原因,我希望在写入和读取时使用Avro模式。
我的理解是,这在Spark之外(或在Spark内手动)是可能的,例如使用AvroParquetWriter和Avro的通用API。但是,我想使用SparkSQL的write()和read()方法(它们与DataFrameWriter和DataFrameReader一起工作),它们与SparkSQL很好地集成在一起(我将编写和读取Dataset)。
我无论如何也想不出如何做到这一点,我想知道这是否可能。SparkSQL拼接格式似乎只支持“压缩”和"mergeSchema“选项--也就是说,没有用于指定替代模式格式或替代模式的选项。换句话说,似乎没有办法通过Avro模式使用SparkSQL应用编程接口来读/写拼图文件。但也许我只是错过了什么?
为了澄清,我也理解,这将基本上只是添加Avro模式到拼花的元数据写入,并将添加一个更多的翻译层读取(Parquet格式-> Avro模式-> SparkSQL内部格式),但将特别允许我添加缺失列的默认值( Avro模式支持,但Parquet模式不)。
此外,我并不是在寻找一种方法来转换Avro到拼图,或者拼图到Avro (而是一种同时使用它们的方式),我也不是在寻找一种在SparkSQL中读/写普通Avro的方法(你可以使用databricks/spark-avro来做到这一点)。
发布于 2017-01-06 02:20:10
我也在做类似的事情。我使用avro模式来写入拼图文件,但是,不要将其读取为avro。但同样的技术也应该适用于read。我不确定这是否是最好的方法,但不管怎样:我有一个具有avro模式的AvroData.avsc。
KafkaUtils.createDirectStream[String,Array[Byte],StringDecoder,DefaultDecoder,Tuple2[String, Array[Byte]]](ssc, kafkaProps, fromOffsets, messageHandler)
kafkaArr.foreachRDD { (rdd,time)
=> { val schema = SchemaConverters.toSqlType(AvroData.getClassSchema).dataType.asInstanceOf[StructType] val ardd = rdd.mapPartitions{itr =>
itr.map { r =>
try {
val cr = avroToListWithAudit(r._2, offsetSaved, loadDate, timeNow.toString)
Row.fromSeq(cr.toArray)
} catch{
case e:Exception => LogHandler.log.error("Exception while converting to Avro" + e.printStackTrace())
System.exit(-1)
Row(0) //This is just to allow compiler to accept. On exception, the application will exit before this point
}
}
}
public static List avroToListWithAudit(byte[] kfkBytes, String kfkOffset, String loaddate, String loadtime ) throws IOException {
AvroData av = getAvroData(kfkBytes);
av.setLoaddate(loaddate);
av.setLoadtime(loadtime);
av.setKafkaOffset(kfkOffset);
return avroToList(av);
}
public static List avroToList(AvroData a) throws UnsupportedEncodingException{
List<Object> l = new ArrayList<>();
for (Schema.Field f : a.getSchema().getFields()) {
String field = f.name().toString();
Object value = a.get(f.name());
if (value == null) {
//System.out.println("Adding null");
l.add("");
}
else {
switch (f.schema().getType().getName()){
case "union"://System.out.println("Adding union");
l.add(value.toString());
break;
default:l.add(value);
break;
}
}
}
return l;
}getAvroData方法需要有从原始字节构造avro对象的代码。我还试图找出一种方法来做到这一点,而不必显式地指定每个属性设置器,但似乎没有这样的方法。
public static AvroData getAvroData (bytes)
{
AvroData av = AvroData.newBuilder().build();
try {
av.setAttr(String.valueOf("xyz"));
.....
}
} 希望能有所帮助
https://stackoverflow.com/questions/41450739
复制相似问题