美文网首页
译:Flink---状态查询 Beta

译:Flink---状态查询 Beta

作者: 雪味伦调 | 来源:发表于2019-02-13 16:15 被阅读0次

    flink 1.7 Google翻译

    注意:可查询状态的客户端API当前处于不断发展的状态,并且不保证所提供接口的稳定性。在即将推出的Flink版本中,客户端可能会发生重大的API更改。
    简而言之,此功能将Flink的托管键控(分区)状态(请参阅使用状态)暴露给外部世界,并允许用户从Flink外部查询作业的状态。对于某些情况,可查询状态消除了对外部系统(例如键值存储)的分布式操作/事务的需要,这通常是实践中的瓶颈。此外,此功能对于调试目的可能特别有用
    注意:查询状态对象时,无需任何同步或复制即可从并发线程访问该对象。这是一种设计选择,因为上述任何一种都会导致增加的作业延迟,我们希望避免这种情况。由于任何状态后端使用Java堆空间,例如MemoryStateBackend或FsStateBackend在检索值时不能与副本一起使用,而是直接引用存储的值,读取 - 修改 - 写入模式是不安全的,并且可能导致可查询状态服务器由于并发修改而失败。 RocksDBStateBackend可以避免这些问题。

    架构


    在展示使用状态查询前,我们先说明状态查询的实体组成部分,这是时分必要的。状态查询由 3部分实体组成

    1. QueryableStateClient,运行在flink集群之外,提交用户查询
    2. QueryableStateClientProxy, 运行在每个TaskManager中(即FLink集群中),它负责接收client的查询,代表client向TaskManager拉取状态信息,并将结果返回client
    3. QueryableStateServer , 它运行在TaskManager,负责服务本地状态存储

    激活状态查询


    在Flink集群开启状态查询,你只需将flink-queryable-state-runtime_2.11-1.7.1.jar从Flink的opt文件夹复制到lib文件夹。否则,状态查询时不可用的。
    通过检查task manager日志中是否有"Started the Queryable State Proxy Server @ ..."判断状态查询是否开启成功

    使状态查询可见


    你已经成功激活状态查询,在使用之前,为使状态对外可见,需要明确一下两点:

    • QueryableStateStream, 一个便利对象,充当接收器并将其传入值作为可查询状态提供
    • stateDescriptor.setQueryable(String queryableStateName)方法,使得由状态描述符表示的键控状态是可查询的

    以下介绍如何使用这两点

    可查询状态流

    在KeyedStream上调用.asQueryableState(stateName,stateDescriptor)会返回一个QueryableStateStream,它将其值提供为可查询状态。根据状态的类型,asQueryableState()方法有以下变体:

    // ValueState
    QueryableStateStream asQueryableState(
        String queryableStateName,
        ValueStateDescriptor stateDescriptor)
    
    // Shortcut for explicit ValueStateDescriptor variant
    QueryableStateStream asQueryableState(String queryableStateName)
    
    // FoldingState
    QueryableStateStream asQueryableState(
        String queryableStateName,
        FoldingStateDescriptor stateDescriptor)
    
    // ReducingState
    QueryableStateStream asQueryableState(
        String queryableStateName,
        ReducingStateDescriptor stateDescriptor)
    

    没有可查询的ListState接收器,因为它会导致不断增长的列表,这些列表可能无法清理,因此最终会消耗太多内存

    返回的QueryableStateStream可以看作是接收器,无法进一步转换。在内部,QueryableStateStream被转换为运算符,该运算符使用所有传入记录来更新可查询状态实例。更新逻辑由asQueryableState调用中提供的StateDescriptor的类型暗示。在如下所示的程序中,键控流的所有记录将用于通过ValueState.update(value)更新状态实例

    stream.keyBy(0).asQueryableState("query-name")
    

    管理监控状态

    通过StateDescriptor.setQueryable(String queryableStateName)查询适当的状态描述符,可以使运算符的托管键控状态(请参阅使用托管键控状态)可查询,如下例所示

    ValueStateDescriptor<Tuple2<Long, Long>> descriptor =
            new ValueStateDescriptor<>(
                    "average", // the state name
                    TypeInformation.of(new TypeHint<Tuple2<Long, Long>>() {})); // type information
    descriptor.setQueryable("query-name"); // queryable state name
    

    queryableStateName参数可以任意选择,仅用于查询。它不必与state自己的名字相同。

    该变体对于哪种类型的状态可以被查询没有限制。这意味着它可以用于任何ValueState,ReduceState,ListState,MapState,AggregatingState和当前不推荐使用的FoldingState。

    状态查询


    现在,你已经设置集群为可查询状态并声明了可查询的状态。是时间去了解如何查询了。
    状态查询需要用到QueryableStateClient 帮助类,它位于flink-queryable-state-client jar包中,你需要显示的在pom中声明引用,并且它与flink core是相互独立的,如下

    <dependency>
      <groupId>org.apache.flink</groupId>
      <artifactId>flink-core</artifactId>
      <version>1.7.1</version>
    </dependency>
    <dependency>
      <groupId>org.apache.flink</groupId>
      <artifactId>flink-queryable-state-client-java_2.11</artifactId>
      <version>1.7.1</version>
    </dependency>
    

    你可以阅读配置Flink程序来了解更多的项目配置。
    QueryableStateClient 向内部代理发送查询请求,代理对象处理查询请求并将结果返回。
    client唯一需要初始化的是提供一个合法的的TaskManager 主机名和代理对象监听的端口号(记着,TaskManager运行着状态可查询的代理)。更多的代理配置及服务端口号请查看Configuration Section

    QueryableStateClient client = new QueryableStateClient(tmHostname, proxyPort);
    

    在客户端准备好的情况下,要查询与类型K的键关联的类型V的状态,您可以使用该方法

    CompletableFuture<S> getKvState(
        JobID jobId,
        String queryableStateName,
        K key,
        TypeInformation<K> keyTypeInfo,
        StateDescriptor<S, V> stateDescriptor)
    

    上面返回一个CompletableFuture,最终保存由具有ID jobID的作业的queryableStateName标识的可查询状态实例的状态值。key是你所感兴趣的状态值,keyTypeInfo 会告知Flink序列化它的方法。最后,stateDescriptor包含有关请求状态的必要信息,即其类型(Value,Reduce等)以及有关如何序列化/反序列化它的必要信息。
    细心的读者会 注意到返回的future包含一个S类值,即一个包含实际值的State对象。这可以是Flink支持的任何状态类型:ValueState,ReduceState,ListState,MapState,AggregatingState和当前不推荐使用的FoldingState。
    这些状态对象不允许修改包含的状态。您可以使用它们来获取状态的实际值,例如使用valueState.get(),或迭代所包含的<K,V>条目,例如使用mapState.entries(),但您无法修改它们。例如,在返回的列表状态上调用add()方法将抛出UnsupportedOperationException

    客户端是异步的,可以由多个线程共享。它需要在未使用时通过QueryableStateClient.shutdown()关闭以释放资源。

    例子

    public class CountWindowAverage extends RichFlatMapFunction<Tuple2<Long, Long>, Tuple2<Long, Long>> {
    
        private transient ValueState<Tuple2<Long, Long>> sum; // a tuple containing the count and the sum
    
        @Override
        public void flatMap(Tuple2<Long, Long> input, Collector<Tuple2<Long, Long>> out) throws Exception {
            Tuple2<Long, Long> currentSum = sum.value();
            currentSum.f0 += 1;
            currentSum.f1 += input.f1;
            sum.update(currentSum);
    
            if (currentSum.f0 >= 2) {
                out.collect(new Tuple2<>(input.f0, currentSum.f1 / currentSum.f0));
                sum.clear();
            }
        }
    
        @Override
        public void open(Configuration config) {
            ValueStateDescriptor<Tuple2<Long, Long>> descriptor =
                    new ValueStateDescriptor<>(
                            "average", // the state name
                            TypeInformation.of(new TypeHint<Tuple2<Long, Long>>() {})); // type information
            descriptor.setQueryable("query-name");
            sum = getRuntimeContext().getState(descriptor);
        }
    }
    

    在作业中使用后,您可以检索作业ID,然后从该运算符查询任何键的当前状态

    // the state descriptor of the state to be fetched.
    ValueStateDescriptor<Tuple2<Long, Long>> descriptor =
            new ValueStateDescriptor<>(
              "average",
              TypeInformation.of(new TypeHint<Tuple2<Long, Long>>() {}));
    
    CompletableFuture<ValueState<Tuple2<Long, Long>>> resultFuture =
            client.getKvState(jobId, "query-name", key, BasicTypeInfo.LONG_TYPE_INFO, descriptor);
    
    // now handle the returned value
    resultFuture.thenAccept(response -> {
            try {
                Tuple2<Long, Long> res = response.get();
            } catch (Exception e) {
                e.printStackTrace();
            }
    });
    

    配置


    QueryableStateOptions 定义了影响状态查询服务和客户端行为的配置参数

    服务状态

    • query.server.ports: 可查询状态服务器的服务器端口范围。如果有多个任务管理器在同一台机器上运行,这对于避免端口冲突很有用。指定的范围可以是:端口:“9123”,一系列端口:“50100-50200”,或范围和/或点列表:“50100-50200,50300-50400,51234”。默认端口为9067。
    • query.server.network-threads: 接收状态服务器传入请求的网络(事件循环)线程数(0 => #slots)
    • query.server.query-threads: 处理/服务状态服务器的传入请求的线程数(0 => #slots)

    代理

    • query.proxy.ports: 可查询状态代理的服务器端口范围。如果有多个任务管理器在同一台机器上运行,这对于避免端口冲突很有用。指定的范围可以是:端口:“9123”,一系列端口:“50100-50200”,或范围和/或点列表:“50100-50200,50300-50400,51234”。默认端口为9069。
    • query.proxy.network-threads: 接收客户端代理的传入请求的网络(事件循环)线程数(0 => #slots)
    • query.proxy.query-threads: 处理/服务客户端代理的传入请求的线程数(0 => #slots)

    限制


    • 可查询状态生命周期与作业的生命周期绑定,例如,任务在启动时注册可查询状态,并在处理时取消注册。在将来的版本中,需要将其解耦以便在任务完成后允许查询,并通过状态复制加速恢复
    • 关于可用KvState的通知是通过一个简单的告诉发生的。在未来,应该通过询问和确认来改进这一点
    • 服务器和客户端会跟踪查询的统计信息。默认情况下,这些功能目前处于禁用状态,因为他们不能再任何地方对外暴露,只要有更好的支持通过Metrics系统发布这些数字,我们就应该启用统计数据

    相关文章

      网友评论

          本文标题:译:Flink---状态查询 Beta

          本文链接:https://www.haomeiwen.com/subject/urldeqtx.html