Как определить собственную логику слияния с помощью Apache Hudi
В Hudi Вы можете настроить класс полезной нагрузки для данной таблицы Hudi по своему усмотрению. Это дает пользователям возможность определить семантику слияния, когда встречаются две версии одной и той же записи. Давайте заглянем под капот, чтобы понять назначение класса полезной нагрузки и то, какие различные способы можно использовать.
Конфигурация: hoodie.datasource.write.payload.class
Примечание: данная информация о классе полезной нагрузки применима ко всем версиям Hudi до 0.15.0. В будущих версиях этот класс может быть неактуален.
Класс Payload
В Hudi есть интерфейс класса Payload, который определяет, как объединяются две версии одной и той же записи.
Выдержка из интерфейса, которая нас интересует:
/**
* This methods lets you write custom merging/combining logic to produce new values as a function of current value on storage and whats contained * in this object. Implementations can leverage properties if required. * <p> * eg: * 1) You are updating counters, you may want to add counts to currentValue and write back updated counts * 2) You may be reading DB redo logs, and merge them with current image for a database row on storage * </p> * * @param currentValue Current value in storage, to merge/combine this payload with * @param schema Schema used for record * @param properties Payload related properties. For example pass the ordering field(s) name to extract from value in storage. * @return new combined/merged value to be written back to storage. EMPTY to skip writing this record. */
Option<IndexedRecord> combineAndGetUpdateValue(IndexedRecord currentValue, Schema schema, Properties properties) throws IOException {
.... }
Внутренне Hudi представляет запись как HoodieRecord, которая состоит из пары HoodieKey и HoodieRecordPayload. HoodieKey, как мы уже видели в предыдущих блогах, представляет собой первичный ключ записи (обычно это путь к разделу и ключ записи). HoodieRecordPayload - это фактические данные, передаваемые пользователем.
Давайте рассмотрим типичный пример. В commit1 Вы вводите 2 записи, а именно {HK1, payload1_1} и {HK2, payload2_1} в commit1(HK -> HoodieKey). В commit2 Вы вводите {HK1, payload1_2} и {HK3, payload3_1}.
Поскольку мы видим обновление для HK1, hudi придется объединить две полезные нагрузки (payload1_1 и payload1_2), чтобы получить окончательный результат для HK1. Вот тут-то и вступает в игру функция combineAndGetUpdateValue(), показанная выше.
По сути, HK1.payload1_2.combineAndGetUpdateValue(HK1.payload1_1) выводит окончательное значение для HK1 в конце commit2.
Варианты реализации полезной нагрузки
Hudi предлагает несколько готовых классов полезной нагрузки для удовлетворения различных потребностей пользователей. Вот лишь некоторые из них: OverwriteWithLatestAvroPayload, DefaultHoodieRecordPayload, AWSDmsAvroPayload, PostgresDebeziumAvroPayload, MySqlDebeziumAvroPayload и т. д.
Но давайте попробуем немного развлечься, создав свой собственный класс Payload.
Класс Payload для суммирования
Допустим, мы хотим найти сумму всех значений столбца в разных версиях одной и той же записи. Допустим, у нас есть таблица сотрудников, в которой есть столбец «зарплата». Посмотрим, как можно найти общую зарплату сотрудника без дополнительных вычислений и без необходимости сканировать все версии записи при выполнении запроса, а также без необходимости хранить все значения зарплаты, которые были найдены до сих пор.
AvroSummationPayload
Давайте определим полезную нагрузку под названием AvroSummationPayload. Мы можем ввести два конфига, чтобы принимать входные данные от пользователя. Один из них будет указывать на столбец, который нужно просуммировать, а второй - на столбец, в котором будет храниться сумма. Мне пришлось ввести этот второй конфиг только для того, чтобы мы могли достичь этого без внесения каких-либо изменений в Hudi, а просто используя определенный пользователем класс полезной нагрузки.
Все, что нужно сделать payload, - это десериализовать новую входящую полезную нагрузку, найти значение в первом конфиге и просуммировать его со значением, хранящимся во втором столбце конфига в уже существующей полезной нагрузке в Hudi.
Сниппет кода
@Override
public Option<IndexedRecord> combineAndGetUpdateValue(IndexedRecord currentValue, Schema schema, Properties properties) throws IOException {
if (recordBytes.length == 0) {
return Option.empty();
} GenericRecord incomingRecord = HoodieAvroUtils.bytesToAvro(recordBytes, schema);
final String sumInputField = properties.getProperty(SUM_INPUT_FIELDS);
final String sumOutputField = properties.getProperty(SUM_OUTPUT_FIELDS);
if (incomingRecord.getSchema().getField(sumOutputField) == null || incomingRecord.getSchema().getField(sumInputField) == null) {
throw new HoodieException("Sum input nor sum output field missing from table schema");
} Long sumInput = (Long) incomingRecord.get(sumInputField);
Long prevTotalSum = (Long) ((GenericRecord) currentValue).get(sumOutputField);
long updatedSum = prevTotalSum + sumInput;
incomingRecord.put(sumOutputField, updatedSum); return Option.of(incomingRecord);
}
Здесь я принял интересующий нас столбец за Long. Но мы должны быть в состоянии приспособиться к любому типу данных. Это всего лишь пробный вариант.
Класс AvroSummationPayload
package org.apache.hudi.common.model;
import org.apache.hudi.avro.HoodieAvroUtils;
import org.apache.hudi.common.util.Option;
import org.apache.hudi.exception.HoodieException;
import org.apache.avro.Schema;
import org.apache.avro.generic.GenericRecord;
import org.apache.avro.generic.IndexedRecord;
import java.io.IOException;
import java.util.Properties;
public class AvroSummationPayload extends BaseAvroPayload
implements HoodieRecordPayload<AvroSummationPayload> {
public static final String SUM_INPUT_FIELDS = "hoodie.summation.payload.input.fields";
public static final String SUM_OUTPUT_FIELDS = "hoodie.summation.payload.output.fields";
public AvroSummationPayload(GenericRecord record, Comparable orderingVal) {
super(record, orderingVal);
} public AvroSummationPayload(Option<GenericRecord> record) {
this(record.isPresent() ? record.get() : null, 0); // natural order
} @Override
public AvroSummationPayload preCombine(AvroSummationPayload oldValue) {
if (oldValue.recordBytes.length == 0) {
// use natural order for delete record
return this;
} if (oldValue.orderingVal.compareTo(orderingVal) > 0) {
// pick the payload with greatest ordering value
return oldValue;
} else {
return this;
} } @Override
public Option<IndexedRecord> combineAndGetUpdateValue(IndexedRecord currentValue, Schema schema) throws IOException {
return getInsertValue(schema);
} @Override
public Option<IndexedRecord> combineAndGetUpdateValue(IndexedRecord currentValue, Schema schema, Properties properties) throws IOException {
if (recordBytes.length == 0) {
return Option.empty();
} GenericRecord incomingRecord = HoodieAvroUtils.bytesToAvro(recordBytes, schema);
final String sumInputField = properties.getProperty(SUM_INPUT_FIELDS);
final String sumOutputField = properties.getProperty(SUM_OUTPUT_FIELDS);
if (incomingRecord.getSchema().getField(sumOutputField) == null || incomingRecord.getSchema().getField(sumInputField) == null) {
throw new HoodieException("Sum input nor sum output field missing from table schema");
} Long sumInput = (Long) incomingRecord.get(sumInputField);
Long prevTotalSum = (Long) ((GenericRecord) currentValue).get(sumOutputField);
long updatedSum = prevTotalSum + sumInput;
incomingRecord.put(sumOutputField, updatedSum); return Option.of(incomingRecord);
} @Override
public Option<IndexedRecord> getInsertValue(Schema schema) throws IOException {
if (recordBytes.length == 0) {
return Option.empty();
} return Option.of(HoodieAvroUtils.bytesToAvro(recordBytes, schema));
} @Override
public Comparable<?> getOrderingValue() {
return this.orderingVal;
} }
Допустим, мы вводим два датафрейма в Hudi, используя вышеуказанный класс payload:
val df = Seq(
(1, "abc", "val1", 10L, 10L), (2, "def", "val1", 20L, 20L), (3, "ghi", "val1", 30L, 30L) ).toDF("id","name", "col1","col2","sum_col2")
val df1 = Seq(
(1, "abc", "val2", 100L, 100L), (2, "def", "val2", 200L, 200L), (3, "ghi", "val2", 300L, 300L) ).toDF("id","name", "col1","col2","sum_col2")
Примечание:
Столбец, который необходимо просуммировать: «col2».
«Sum_col2» добавляется как новый столбец со значением, аналогичным „col2“. Этот столбец будет содержать общую сумму для col2 после завершения записи.
При записи в hudi нам нужно установить следующие дополнительные параметры записи.
DataSourceWriteOptions.RECORDKEY_FIELD.key -> "id", DataSourceWriteOptions.PRECOMBINE_FIELD.key -> "id", DataSourceWriteOptions.PARTITIONPATH_FIELD.key -> "", "hoodie.datasource.write.payload.class" -> "org.apache.hudi.common.model.AvroSummationPayload", "hoodie.summation.payload.input.fields" -> "col2", "hoodie.summation.payload.output.fields" -> "sum_col2",
Примечание: я задал поле precombine таким же, как и ключ записи, поскольку это просто пример.
Когда мы читаем снимок Hudi после поглощения двух вышеуказанных партий, вывод выглядит так:
select id, sum_col2 from hudi_tbl + - -+ - - - - - - - + |id |sum_output_col| + - -+ - - - - - - - + |1 |110 | |3 |330 | |2 |220 | + - -+ - - - - - - - +
Запись с идентификатором 1: 10 + 100 = 110 Запись с идентификатором 2: 20 + 200 = 220 Запись с идентификатором 3: 30 + 300 = 330
Payload для нахождения среднего значения для заданного столбца по ключу записи
Можно ли попытаться определить полезную нагрузку для поиска среднего значения по заданному столбцу.
Суммирование может быть простым, поскольку нам нужно просто продолжать добавлять значения одного и того же столбца, не требуя никакой дополнительной информации. Но для среднего значения нам также нужно хранить количество записей вместе с итоговым значением avg.
Как и ожидалось, этот payload очень похож на AvroSummationPayload, с небольшими улучшениями для хранения общего количества записей в дополнение к среднему значению.
Gist payload подразумевает нахождение среднего значения:
@Override
public Option<IndexedRecord> combineAndGetUpdateValue(IndexedRecord currentValue, Schema schema, Properties properties) throws IOException {
if (recordBytes.length == 0) {
return Option.empty();
} GenericRecord incomingRecord = HoodieAvroUtils.bytesToAvro(recordBytes, schema);
final String avgInputField = properties.getProperty(AVG_INPUT_FIELDS);
final String avgOutputFieldValue = properties.getProperty(AVG_OUTPUT_FIELDS_VALUE);
final String avgOutputFieldCount = properties.getProperty(AVG_OUTPUT_FIELDS_COUNT);
if (incomingRecord.getSchema().getField(avgInputField) == null || incomingRecord.getSchema().getField(avgOutputFieldValue) == null) {
throw new HoodieException("Sum input nor sum output field missing from table schema");
} Long newInput = (Long) incomingRecord.get(avgInputField);
Double prevTotalAverage = (Double) ((GenericRecord) currentValue).get(avgOutputFieldValue);
Long prevTotalValues = (Long) ((GenericRecord) currentValue).get(avgOutputFieldCount);
double updatedSum = prevTotalAverage * prevTotalValues;
updatedSum += newInput; incomingRecord.put(avgOutputFieldValue, updatedSum/(prevTotalValues + 1));
incomingRecord.put(avgOutputFieldCount, prevTotalValues + 1);
return Option.of(incomingRecord);
}j
Столбец AVG_OUTPUT_FIELDS_VALUE содержит среднее значение для всех значений AVG_INPUT_FIELDS, наблюдавшихся до сих пор, когда мы выполняли моментальное чтение интересующей нас таблицы.
Вот Gist ccылка на оба payload (сумма и среднее значение) для фолков.
Мощная гибкость
Собственная реализация payload позволяет при необходимости внедрять пользовательскую бизнес-логику. Две реализации payload, которые мы рассмотрели выше, также могут снизить затраты на вычисления. Если бы не эти пользовательские реализации payload, и если бы нужно было вычислить сумму или среднее значение по всем версиям записи, нам, возможно, пришлось бы присоединиться к целевой таблице w/, чтобы найти сумму или среднее значение, а затем обновить целевую таблицу sink (золотой слой). В качестве альтернативы может потребоваться сканирование всей серебряной таблицы для повторного вычисления таблицы золотого слоя. Вместо этого мы могли бы выполнять инкрементное занесение данных, сохраняя ту же семантику целевой таблицы. Это может привести к огромной экономии средств при правильном применении. Кстати, приведенные выше две реализации payload были разработаны специально для этого блога. Они не доступны ни в одном из официальных релизов hudi.
Примечание: Пользователи, использующие spark-sql, могут использовать любые выражения, как обычно, и им не нужно использовать эти пользовательские payload. Но это может быть полезно для пользователей, использующих spark datasource writes или HoodieStreamer или HoodieSparkStreaming для ввода данных в Hudi.




