Противоречие Flink Async — Sharp Async I/O

Java

Таблица измерений JOIN — бизнес-сценарий, который нельзя обойти

В процессе обработки потока Flink часто необходимо взаимодействовать с внешней системой и использовать таблицу измерений для заполнения полей в таблице фактов.

Например: в сценарии электронной коммерции skuid продукта требуется для сопоставления некоторых атрибутов продукта, таких как отрасль, к которой принадлежит продукт, производитель продукта и некоторая информация о производителе; в логистике сценарий, если вы знаете идентификатор пакета, вам необходимо связать идентификатор пакета.Отраслевые атрибуты, информация о доставке, информация о получении и т. д.

По умолчанию в MapFunction от Flink один параллель может взаимодействовать только синхронно: отправить запрос во внешнее хранилище, заблокировать ввод-вывод, дождаться возврата запроса, а затем продолжить отправку следующего запроса. Такой способ синхронного взаимодействия часто занимает много времени в сети. Чтобы повысить эффективность обработки, можно увеличить параллелизм MapFunction, но увеличение параллелизма означает больше ресурсов, что является не очень хорошим решением.

Асинхронные неблокирующие запросы ввода/вывода

Flink представил асинхронный ввод-вывод в версии 1.2. В асинхронном режиме операции ввода-вывода являются асинхронными. Один параллельный запрос может отправлять несколько запросов подряд, и тот запрос, который возвращается первым, будет обработан первым, поэтому блокировка ожидания не требуется между последовательными запросами. что значительно повышает эффективность потоковой обработки.

Асинхронный ввод-вывод — очень популярная функция, предоставленная сообществу Alibaba, которая решает проблему сетевых задержек, которые становятся узким местом системы при взаимодействии с внешними системами.

file

Длинная коричневая полоса на рисунке представляет время ожидания, и можно обнаружить, что время ожидания в сети сильно снижает пропускную способность и задержку. Чтобы решить проблему синхронного доступа, асинхронный режим может обрабатывать несколько запросов и ответов одновременно. То есть вы можете непрерывно отправлять запросы от пользователей a, b, c и т. д. в базу данных, и в то же время будет обработан тот ответ, который будет возвращен первым, так что нет необходимости блокировать и ждать между последовательные запросы, как показано на рисунке выше, показанном справа. Именно так работает асинхронный ввод-вывод.

Подробный принцип приведен в первой ссылке в конце статьи, которой поделилась Alibaba Yunxie.

Простой пример выглядит следующим образом:

public class AsyncIOFunctionTest {
    public static void main(String[] args) throws Exception {

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
        env.setParallelism(1);

        Properties p = new Properties();
        p.setProperty("bootstrap.servers", "localhost:9092");

        DataStreamSource<String> ds = env.addSource(new FlinkKafkaConsumer010<String>("order", new SimpleStringSchema(), p));
        ds.print();

        SingleOutputStreamOperator<Order> order = ds
                .map(new MapFunction<String, Order>() {
                    @Override
                    public Order map(String value) throws Exception {
                        return new Gson().fromJson(value, Order.class);
                    }
                })
                .assignTimestampsAndWatermarks(new AscendingTimestampExtractor<Order>() {
                    @Override
                    public long extractAscendingTimestamp(Order element) {
                        try {
                            return element.getOrderTime();
                        } catch (Exception e) {
                            e.printStackTrace();
                        }
                        return 0;
                    }
                })
                .keyBy(new KeySelector<Order, String>() {
                    @Override
                    public String getKey(Order value) throws Exception {
                        return value.getUserId();
                    }
                })
                .window(TumblingEventTimeWindows.of(Time.minutes(10)))
                .maxBy("orderTime");

        SingleOutputStreamOperator<Tuple7<String, String, Integer, String, String, String, Long>> operator = AsyncDataStream
                .unorderedWait(order, new RichAsyncFunction<Order, Tuple7<String, String, Integer, String, String, String, Long>>() {

                    private Connection connection;

                    @Override
                    public void open(Configuration parameters) throws Exception {
                        super.open(parameters);
                        Class.forName("com.mysql.jdbc.Driver");
                        connection = DriverManager.getConnection("url", "user", "pwd");
                        connection.setAutoCommit(false);
                    }

                    @Override
                    public void asyncInvoke(Order input, ResultFuture<Tuple7<String, String, Integer, String, String, String, Long>> resultFuture) throws Exception {
                        List<Tuple7<String, String, Integer, String, String, String, Long>> list = new ArrayList<>();
                        // 在 asyncInvoke 方法中异步查询数据库
                        String userId = input.getUserId();
                        Statement statement = connection.createStatement();
                        ResultSet resultSet = statement.executeQuery("select name,age,sex from user where userid=" + userId);
                        if (resultSet != null && resultSet.next()) {
                            String name = resultSet.getString("name");
                            int age = resultSet.getInt("age");
                            String sex = resultSet.getString("sex");
                            Tuple7<String, String, Integer, String, String, String, Long> res = Tuple7.of(userId, name, age, sex, input.getOrderId(), input.getPrice(), input.getOrderTime());
                            list.add(res);
                        }

                        // 将数据搜集
                        resultFuture.complete(list);
                    }

                    @Override
                    public void close() throws Exception {
                        super.close();
                        if (connection != null) {
                            connection.close();
                        }
                    }
                }, 5000, TimeUnit.MILLISECONDS,100);

        operator.print();


        env.execute("AsyncIOFunctionTest");
    }
}

В приведенном выше коде исходный поток заказов поступает из Kafka, а пользовательская информация о заказе извлекается из связанной таблицы измерений. Как видно из приведенного выше примера, мы создаем объект соединения в open(), закрываем соединение в методе close() и в методе asyncInvoke() RichAsyncFunction напрямую запрашиваем операцию базы данных и возвращаем данные. Такой простой асинхронный запрос выполнен.

Принцип и базовое использование асинхронного ввода-вывода

Проще говоря, API, соответствующий Flink, использующему асинхронный ввод-вывод, представляет собой абстрактный класс RichAsyncFunction, который реализует методы open (инициализация), asyncInvoke (асинхронный вызов данных) и close (некоторые операции для остановки) в этом абстрактном классе. Самое главное — реализовать методы в asyncInvoke.

Давайте сначала рассмотрим метод шаблона, использующий асинхронный ввод-вывод:


// This example implements the asynchronous request and callback with Futures that have the
// interface of Java 8's futures (which is the same one followed by Flink's Future)

/**
 * An implementation of the 'AsyncFunction' that sends requests and sets the callback.
 */
class AsyncDatabaseRequest extends RichAsyncFunction<String, Tuple2<String, String>> {

    /** The database specific client that can issue concurrent requests with callbacks */
    private transient DatabaseClient client;

    @Override
    public void open(Configuration parameters) throws Exception {
        client = new DatabaseClient(host, post, credentials);
    }

    @Override
    public void close() throws Exception {
        client.close();
    }

    @Override
    public void asyncInvoke(String key, final ResultFuture<Tuple2<String, String>> resultFuture) throws Exception {

        // issue the asynchronous request, receive a future for result
        final Future<String> result = client.query(key);

        // set the callback to be executed once the request by the client is complete
        // the callback simply forwards the result to the result future
        CompletableFuture.supplyAsync(new Supplier<String>() {

            @Override
            public String get() {
                try {
                    return result.get();
                } catch (InterruptedException | ExecutionException e) {
                    // Normally handled explicitly.
                    return null;
                }
            }
        }).thenAccept( (String dbResult) -> {
            resultFuture.complete(Collections.singleton(new Tuple2<>(key, dbResult)));
        });
    }
}

// create the original stream
DataStream<String> stream = ...;

// apply the async I/O transformation
DataStream<Tuple2<String, String>> resultStream =
    AsyncDataStream.unorderedWait(stream, new AsyncDatabaseRequest(), 1000, TimeUnit.MILLISECONDS, 100);

Предположим, нам нужно сделать асинхронные запросы к другим базам данных в сценарии, тогда для реализации операции с базой данных через асинхронный ввод-вывод требуется три шага: 1. Реализовать AsyncFunction, используемую для распределения запроса 2. Получить обратный вызов результата операции и поместить Отправить в AsyncCollector 3. Преобразовать асинхронные операции ввода-вывода в DataStream

    其中的两个重要的参数:

TimeouttimeoutОпределяет, как долго будет отбрасываться асинхронная операция. Этот параметр используется для предотвращения мертвых или неудачных запросов.CapacityЭтот параметр определяет, сколько асинхронных запросов может обрабатываться одновременно. Хотя подход с асинхронным вводом-выводом обеспечивает лучшую пропускную способность, операторы по-прежнему могут быть узким местом для потоковых приложений. Превышение ограничения на количество одновременных запросов создает обратное давление.

Несколько замечаний:

  • Для использования асинхронного ввода-вывода требуется, чтобы во внешнем хранилище были клиенты, поддерживающие асинхронные запросы.
  • Используйте асинхронный ввод-вывод, наследуйте RichAsyncFunction (абстрактный класс интерфейса AsyncFunction) и перепишите или реализуйте три метода open (установление соединения), close (закрытие соединения) и asyncInvoke (асинхронный вызов).
  • Использование асинхронного ввода-вывода, предпочтительно в сочетании с кешем, может уменьшить количество запросов к внешнему хранилищу и повысить эффективность.
  • Асинхронный ввод-вывод предоставляет параметр Timeout для управления максимальным временем ожидания запросов. По умолчанию, когда время ожидания запроса асинхронного ввода-вывода истекает, создается исключение, и задание перезапускается или останавливается. Если вы хотите обрабатывать тайм-ауты, вы можете переопределить метод AsyncFunction#timeout.
  • Асинхронный ввод-вывод предоставляет параметр «Емкость» для управления количеством одновременных запросов.После того, как емкость будет исчерпана, будет запущен механизм обратного давления, чтобы подавить прием восходящих данных.
  • Выход асинхронного ввода-вывода обеспечивает как неупорядоченный, так и последовательный режимы.
        乱序, 用AsyncDataStream.unorderedWait(...) API,每个并行的输出顺序和输入顺序可能不一致。
        顺序, 用AsyncDataStream.orderedWait(...) API,每个并行的输出顺序和输入顺序一致。为保证顺序,需要在输出的Buffer中排序,该方式效率会低一些。

Оптимизация для Flink 1.9

Благодаря недавно включенным функциям, связанным с Blink, очень просто реализовать функцию таблицы измерений в Flink 1.9. Если вы хотите использовать эту функцию, вам нужно самостоятельно представить планировщик Blink.

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-table-planner-blink_${scala.binary.version}</artifactId>
    <version>${flink.version}</version>
</dependency>

Далее нам нужно только настроить реализацию интерфейса LookupableTableSource и реализовать в нем методы.Проанализируем код LookupableTableSource:

public interface LookupableTableSource<T> extends TableSource<T> {
     TableFunction<T> getLookupFunction(String[] lookupKeys);
     AsyncTableFunction<T> getAsyncLookupFunction(String[] lookupKeys);
     boolean isAsyncEnabled();
}

Эти три метода:

  • isAsyncEnabledМетод в основном указывает, поддерживает ли таблица асинхронный доступ к внешним источникам данных для получения данных.Когда он возвращает true, после регистрации в TableEnvironment он возвращает асинхронную функцию для вызова при ее использовании, а когда возвращает false, сделать функцию синхронного доступа.
  • getLookupFunctionМетод возвращает функцию, которая синхронно обращается к внешней системе данных. Что это значит? Вам нужно запросить внешнюю базу данных через ключ, и вам нужно дождаться возврата данных, прежде чем продолжить обработку данных, что повлияет на пропускную способность скорость обработки системы.
  • getAsyncLookupFunctionМетод заключается в возврате асинхронной функции для асинхронного доступа к внешней системе данных для получения данных, что может значительно повысить пропускную способность системы.

Независимо от функции синхронного доступа, getAsyncLookupFunction вернет функцию, которая асинхронно обращается к внешним источникам данных.Если вы хотите использовать асинхронные функции, предпосылка заключается в том, что метод isAsyncEnabled LookupableTableSource возвращает значение true для его использования. Используйте асинхронные функции для доступа к внешним системам данных. Как правило, внешние системы имеют клиентов асинхронного доступа. Если нет, вы можете использовать пулы потоков для асинхронного доступа к внешним системам. Например:

public class MyAsyncLookupFunction extends AsyncTableFunction<Row> {
    private transient RedisAsyncCommands<String, String> async;
    @Override
    public void open(FunctionContext context) throws Exception {
        RedisClient redisClient = RedisClient.create("redis://127.0.0.1:6379");
        StatefulRedisConnection<String, String> connection = redisClient.connect();
        async = connection.async();
    }
    public void eval(CompletableFuture<Collection<Row>> future, Object... params) {
        redisFuture.thenAccept(new Consumer<String>() {
            @Override
            public void accept(String value) {
                future.complete(Collections.singletonList(Row.of(key, value)));
            }
        });
    }
}

Полный пример выглядит следующим образом:

Основной метод:


import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.api.java.typeutils.RowTypeInfo;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer011;
import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.java.StreamTableEnvironment;
import org.apache.flink.types.Row;
import org.junit.Test;
 
import java.util.Properties;
 
public class LookUpAsyncTest {
 
    @Test
    public void test() throws Exception {
        LookUpAsyncTest.main(new String[]{});
    }
 
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        //env.setParallelism(1);
        EnvironmentSettings settings = EnvironmentSettings.newInstance().useBlinkPlanner().inStreamingMode().build();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env, settings);
 
        final ParameterTool params = ParameterTool.fromArgs(args);
        String fileName = params.get("f");
        DataStream<String> source = env.readTextFile("hdfs://172.16.44.28:8020" + fileName, "UTF-8");
 
        TypeInformation[] types = new TypeInformation[]{Types.STRING, Types.STRING, Types.LONG};
        String[] fields = new String[]{"id", "user_click", "time"};
        RowTypeInfo typeInformation = new RowTypeInfo(types, fields);
 
        DataStream<Row> stream = source.map(new MapFunction<String, Row>() {
            private static final long serialVersionUID = 2349572543469673349L;
 
            @Override
            public Row map(String s) {
                String[] split = s.split(",");
                Row row = new Row(split.length);
                for (int i = 0; i < split.length; i++) {
                            
                    Object value = split[i];
                    if (types[i].equals(Types.STRING)) {
                        value = split[i];
                    }
                    if (types[i].equals(Types.LONG)) {
                        value = Long.valueOf(split[i]);
                    }
                    row.setField(i, value);
                }
                return row;
            }
        }).returns(typeInformation);
 
        tableEnv.registerDataStream("user_click_name", stream, String.join(",", typeInformation.getFieldNames()) + ",proctime.proctime");
 
        RedisAsyncLookupTableSource tableSource = RedisAsyncLookupTableSource.Builder.newBuilder()
                .withFieldNames(new String[]{"id", "name"})
                .withFieldTypes(new TypeInformation[]{Types.STRING, Types.STRING})
                .build();
        tableEnv.registerTableSource("info", tableSource);
 
        String sql = "select t1.id,t1.user_click,t2.name" +
                " from user_click_name as t1" +
                " join info FOR SYSTEM_TIME AS OF t1.proctime as t2" +
                " on t1.id = t2.id";
 
        Table table = tableEnv.sqlQuery(sql);
 
        DataStream<Row> result = tableEnv.toAppendStream(table, Row.class);
 
        DataStream<String> printStream = result.map(new MapFunction<Row, String>() {
            @Override
            public String map(Row value) throws Exception {
                return value.toString();
            }
        });
 
        Properties properties = new Properties();
        properties.setProperty("bootstrap.servers", "127.0.0.1:9094");
        FlinkKafkaProducer011<String> kafkaProducer = new FlinkKafkaProducer011<>(
                "user_click_name",  
                new SimpleStringSchema(),
                properties);
        printStream.addSink(kafkaProducer);
 
        tableEnv.execute(Thread.currentThread().getStackTrace()[1].getClassName());
    }
}

Метод RedisAsyncLookupTableSource:

import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.typeutils.RowTypeInfo;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.TableSchema;
import org.apache.flink.table.functions.AsyncTableFunction;
import org.apache.flink.table.functions.TableFunction;
import org.apache.flink.table.sources.LookupableTableSource;
import org.apache.flink.table.sources.StreamTableSource;
import org.apache.flink.table.types.DataType;
import org.apache.flink.table.types.utils.TypeConversions;
import org.apache.flink.types.Row;
 
public class RedisAsyncLookupTableSource implements StreamTableSource<Row>, LookupableTableSource<Row> {
 
    private final String[] fieldNames;
    private final TypeInformation[] fieldTypes;
 
    public RedisAsyncLookupTableSource(String[] fieldNames, TypeInformation[] fieldTypes) {
       this.fieldNames = fieldNames;
        this.fieldTypes = fieldTypes;
    }
 
    //同步方法
    @Override
    public TableFunction<Row> getLookupFunction(String[] strings) {
        return null;
    }
 
    //异步方法
    @Override
    public AsyncTableFunction<Row> getAsyncLookupFunction(String[] strings) {
        return MyAsyncLookupFunction.Builder.getBuilder()
                .withFieldNames(fieldNames)
                .withFieldTypes(fieldTypes)
                .build();
    }
 
    //开启异步
    @Override
    public boolean isAsyncEnabled() {
        return true;
    }
 
    @Override
    public DataType getProducedDataType() {
        return TypeConversions.fromLegacyInfoToDataType(new RowTypeInfo(fieldTypes, fieldNames));
    }
 
    @Override
    public TableSchema getTableSchema() {
        return TableSchema.builder()
                .fields(fieldNames, TypeConversions.fromLegacyInfoToDataType(fieldTypes))
                .build();
    }
 
    @Override
    public DataStream<Row> getDataStream(StreamExecutionEnvironment environment) {
        throw new UnsupportedOperationException("do not support getDataStream");
    }
 
    public static final class Builder {
        private String[] fieldNames;
        private TypeInformation[] fieldTypes;
 
        private Builder() {
        }
 
        public static Builder newBuilder() {
            return new Builder();
        }
 
        public Builder withFieldNames(String[] fieldNames) {
            this.fieldNames = fieldNames;
            return this;
        }
 
        public Builder withFieldTypes(TypeInformation[] fieldTypes) {
            this.fieldTypes = fieldTypes;
            return this;
        }
 
        public RedisAsyncLookupTableSource build() {
            return new RedisAsyncLookupTableSource(fieldNames, fieldTypes);
        }
    }
}

MyAsyncLookupFunction


import io.lettuce.core.RedisClient;
import io.lettuce.core.RedisFuture;
import io.lettuce.core.api.StatefulRedisConnection;
import io.lettuce.core.api.async.RedisAsyncCommands;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.typeutils.RowTypeInfo;
import org.apache.flink.table.functions.AsyncTableFunction;
import org.apache.flink.table.functions.FunctionContext;
import org.apache.flink.types.Row;
 
import java.util.Collection;
import java.util.Collections;
import java.util.concurrent.CompletableFuture;
import java.util.function.Consumer;
 
public class MyAsyncLookupFunction extends AsyncTableFunction<Row> {
 
    private final String[] fieldNames;
    private final TypeInformation[] fieldTypes;
 
    private transient RedisAsyncCommands<String, String> async;
 
    public MyAsyncLookupFunction(String[] fieldNames, TypeInformation[] fieldTypes) {
        this.fieldNames = fieldNames;
        this.fieldTypes = fieldTypes;
    }
 
    @Override
    public void open(FunctionContext context) {
        //配置redis异步连接
        RedisClient redisClient = RedisClient.create("redis://127.0.0.1:6379");
        StatefulRedisConnection<String, String> connection = redisClient.connect();
        async = connection.async();
    }
 
    //每一条流数据都会调用此方法进行join
    public void eval(CompletableFuture<Collection<Row>> future, Object... paramas) {
        //表名、主键名、主键值、列名
        String[] info = {"userInfo", "userId", paramas[0].toString(), "userName"};
        String key = String.join(":", info);
        RedisFuture<String> redisFuture = async.get(key);
 
        redisFuture.thenAccept(new Consumer<String>() {
            @Override
            public void accept(String value) {
                future.complete(Collections.singletonList(Row.of(key, value)));
                //todo
//                BinaryRow row = new BinaryRow(2);
            }
        });
    }
 
    @Override
    public TypeInformation<Row> getResultType() {
        return new RowTypeInfo(fieldTypes, fieldNames);
    }
 
    public static final class Builder {
        private String[] fieldNames;
        private TypeInformation[] fieldTypes;
 
        private Builder() {
        }
 
        public static Builder getBuilder() {
            return new Builder();
        }
 
        public Builder withFieldNames(String[] fieldNames) {
            this.fieldNames = fieldNames;
            return this;
        }
 
        public Builder withFieldTypes(TypeInformation[] fieldTypes) {
            this.fieldTypes = fieldTypes;
            return this;
        }
 
        public MyAsyncLookupFunction build() {
            return new MyAsyncLookupFunction(fieldNames, fieldTypes);
        }
    }
}

Несколько моментов, на которые стоит обратить внимание:

1. Внешний источник данных должен быть асинхронным клиентом: если он потокобезопасен (используется несколькими клиентами), вы можете инициализировать его один раз, не добавляя ключевое слово transient. В противном случае вам нужно добавить переходный процесс, не инициализируйте его, а в методе open инициализируйте по одному для каждого экземпляра Task.

2. В методе eval есть еще один CompletableFuture, по завершении асинхронного доступа необходимо вызвать его метод для обработки. Например, в приведенном выше примере:

redisFuture.thenAccept(new Consumer<String>() {
            @Override
            public void accept(String value) {
                future.complete(Collections.singletonList(Row.of(key, value)));
            }
        });

3. Хотя в сообществе предусмотрена функция асинхронного связывания таблиц измерений, на самом деле связывание таблиц измерений с внешними системами все равно станет узким местом системы при больших объемах данных, поэтому обычно к синхронным и асинхронным функциям добавляем кэши. С учетом таких факторов, как параллелизм, простота использования, обновления в реальном времени и несколько версий, Hbase является наиболее идеальной внешней таблицей измерений.

Справочная статья: http://wuchong.me/blog/2017/05/17/flink-internals-async-io/#https://www.jianshu.com/p/d8f99d94b761https://cwiki.apache.org/ слияние/страницы/viewpage.action?pageId=65870673https://www.jianshu.com/p/7ce84f978ae0

Подпишитесь на мой официальный аккаунт и ответьте на [JAVAPDF] в фоновом режиме, чтобы получить 200 страниц тестовых вопросов!Большие данные, на которые обращают внимание 50 000 человек, — это дорога к Богу, почему бы вам не прийти и не узнать об этом?50 000 человек обращают внимание на то, как большие данные становятся богом, разве вы не хотите узнать об этом?50 000 человек обращают внимание на то, как большие данные становятся богом, вы уверены, что действительно не хотите прийти и узнать об этом?

Добро пожаловать, чтобы обратить внимание«Дорога к большим данным становится Богом»

大数据技术与架构