Как масштабировать последовательный источник в потоке данных/луче?JAVA

Программисты JAVA общаются здесь
Ответить
Anonymous
 Как масштабировать последовательный источник в потоке данных/луче?

Сообщение Anonymous »

Вопрос заключается в последовательном отправке источника в потоке данных?
Например, реализация bigqueryio go очень проста. Он считывает результаты запроса в цикле и последовательно выдает записи. Поэтому поток данных не масштабирует количество воркеров и обнаруживает отстающего.
Псевдокод такого пользовательского DoFn в Java:

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

@ProcessElement
public void processElement(@Element byte[] input, OutputReceiver output) throws Exception {
// ..
result = bigquery.query(queryConfig);
for (FieldValueList row : result.iterateAll()) {
output.output(mapRowToType(row));
}
}
Последующие шаги, использующие исходный код в ParDo, никогда не масштабируются. Документы намекают, что это может произойти, если пользовательский источник не реализует прогресс (а это не так).
Пример преобразования вышеуказанного источника в Java:

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

var bigqueryRows = pipeline.apply("ReadFromBigQuery", BigQueryIO.read(//...));
var mutations = bigqueryRows.apply("ProcessRows", ParDo.of(/**/ {
@ProcessElement
public void processElement(@Element ItemRow row, OutputReceiver out) {

String rowKey = row.getId();
BigtableIO.Write.Mutation mutation = BigtableIO.Write.Mutation.create(rowKey);
// ...
out.output(mutation);
}
}));

Как быть с таким последовательным источником, если нет другого варианта использования данных.
Есть ли стратегия, которую можно позже добавить в конвейер, чтобы он начал распараллеливаться обработка нескольким работникам?
*ПРИМЕЧАНИЕ:
Хотя моя конкретная проблема связана с go sdk и bigqueryio, мой вопрос касается такого сценария в целом. Не стесняйтесь давать общий ответ на любом знакомом SDK языка луча.

Подробнее здесь: https://stackoverflow.com/questions/792 ... aflow-beam
Ответить

Быстрый ответ

Изменение регистра текста: 
Смайлики
:) :( :oops: :roll: :wink: :muza: :clever: :sorry: :angel: :read: *x)
Ещё смайлики…
   
К этому ответу прикреплено по крайней мере одно вложение.

Если вы не хотите добавлять вложения, оставьте поля пустыми.

Максимально разрешённый размер вложения: 15 МБ.

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