【发布时间】:2017-08-20 13:14:59
【问题描述】:
这里有一个类似的问题:How to add a schema to a Dataset in Spark?
但是我面临的问题是我已经预定义了Dataset<Obj1>,并且我想定义一个模式来匹配它的数据成员。最终目标是能够连接两个 java 对象。
示例代码:
Dataset<Row> rowDataset = spark.getSpark().sqlContext().createDataFrame(rowRDD, schema).toDF();
Dataset<MyObj> objResult = rowDataset.map((MapFunction<Row, MyObj>) row ->
new MyObj(
row.getInt(row.fieldIndex("field1")),
row.isNullAt(row.fieldIndex("field2")) ? "" : row.getString(row.fieldIndex("field2")),
row.isNullAt(row.fieldIndex("field3")) ? "" : row.getString(row.fieldIndex("field3")),
row.isNullAt(row.fieldIndex("field4")) ? "" : row.getString(row.fieldIndex("field4"))
), Encoders.javaSerialization(MyObj.class));
如果我正在打印行数据集的架构,我会按预期获得架构:
rowDataset.printSchema();
root
|-- field1: integer (nullable = false)
|-- field2: string (nullable = false)
|-- field3: string (nullable = false)
|-- field4: string (nullable = false)
如果我正在打印对象数据集,我将丢失实际架构
objResult.printSchema();
root
|-- value: binary (nullable = true)
问题是如何为Dataset<MyObj> 应用架构?
【问题讨论】:
-
请提供一个最小、完整和可验证的示例 - stackoverflow.com/help/mcve
-
关于您的问题的代码可能会帮助我们推荐一些东西。
-
@squid,我提供了一个代码sn-p
标签: java apache-spark spark-dataframe apache-spark-dataset