Добавление новых строк в раздел Spark при использовании forEachPartitionJAVA

Программисты JAVA общаются здесь
Anonymous
Добавление новых строк в раздел Spark при использовании forEachPartition

Сообщение Anonymous »

Я пытаюсь добавить новую строку в каждый раздел в моем задании Spark. Для этого я использую следующий код:

Код: Выделить всё

StructType rowType = new StructType();
rowType.add(DataTypes.createStructField("value", DataTypes.StringType, true));
sourceDataset = sourceDataset.mapPartitions(new MapPartitionsFunction() {
List rows = new ArrayList();
@Override
public Iterator call(Iterator rowIterator) throws Exception {

int counter =0;
while (rowIterator.hasNext()) {
String dataRow = rowIterator.next();
rows.add(dataRow);
counter++;
}

JsonObject jsonObject = func.get();
String[] values = new String[1];
values[0] = jsonObject.toString();
rows.add(values[0]);
return rows.iterator();
}
}
,RowEncoder.apply(rowType)
);
И я получаю следующее исключение.

Код: Выделить всё

org.apache.spark.sql.AnalysisException: Try to map struct to Tuple1, but failed as the number of fields does not line up
В исходном наборе данных у нас есть «значение» строкового типа и это Json в формате String.
Я пробовал изменить MapPartitionsFunction в MapPartitionsFunction и все равно получаю то же исключение.
Спасибо
Сатиш

Подробнее здесь: https://stackoverflow.com/questions/787 ... hpartition

Вернуться в «JAVA»