【问题标题】:Access flink state inside processBroadcastElement function在 processBroadcastElement 函数中访问 flink 状态
【发布时间】:2021-04-13 02:29:58
【问题描述】:

打算在processBroadcastElement()函数内部做一些状态管理。

final val actvTagsMapValue = new MapStateDescriptor[String, List[String]]("actvTagsMapValue", classOf[String], classOf[List[String]])

override def processBroadcastElement(...): Unit {
    val actvTagMap = getRuntimeContext.getMapState(actvTagsMapValue)
    val st = actvTagMap.entries() // this line produce an error
}

访问状态时出现以下错误

229797 [LabelShlfEvents -> Sink: Print to Std. Out (1/1)] WARN  
o.a.flink.runtime.taskmanager.Task - LabelShlfEvents -> Sink: Print to Std. Out (1/1) 
(d3154841fd8bd4cabc00e0145ac37ed8) switched from RUNNING to FAILED. 
java.lang.NullPointerException: No key set. This method should not be called outside of a keyed context.
at org.apache.flink.util.Preconditions.checkNotNull(Preconditions.java:75)

我不能这样做吗?

【问题讨论】:

    标签: apache-flink


    【解决方案1】:

    这不起作用,Dominik 在his answer 中解释了原因。

    您可以在processBroadcastElement 中执行的操作是访问/修改/删除所有键的键控状态,方法是使用applyToKeyedStateKeyedStateFunction。但是,您必须注意在所有并行实例中的行为具有确定性。否则,在恢复或重新缩放后,您可能会出现不一致。

    这是一个示例,它在接收到任何广播消息时为每个键发出 ValueState 的值。

    public static class DumpFunction
            extends KeyedBroadcastProcessFunction<Long, TaxiRide, String, TaxiRide> {
        private ValueStateDescriptor<TaxiRide> taxiDesc;
        private ValueState<TaxiRide> taxiState;
        
        @Override
        public void open(Configuration config) {
            taxiDesc = new ValueStateDescriptor<>("ride", TaxiRide.class);
            taxiState = getRuntimeContext().getState(taxiDescriptor);
        }
        
        @Override
        public void processElement(TaxiRide ride, ReadOnlyContext ctx, 
                Collector< TaxiRide> out) throws Exception {
            taxiState.update(ride);
        }
                
        @Override
        public void processBroadcastElement(String msg, Context ctx, Collector<TaxiRide> out) {
            ctx.applyToKeyedState(taxiDesc, new KeyedStateFunction<Long, ValueState<TaxiRide>>() {
                @Override
                public void process(Long taxiId, ValueState<TaxiRide> taxiState) throws Exception {
                    out.collect(taxiState.value());
                }
            });
        }
    }
    

    【讨论】:

      【解决方案2】:

      你不能这样做。原因很简单,因为MapState(还有ValueStateListState 和更多描述的here)是一种称为键控状态的状态。此状态被分区并作用于当前元素的输入键。

      广播元素没有以任何方式进行键控或分区,因此这些元素没有附加KeyedContext。当你尝试访问processBroadcastElement 中的状态时,Flink 不知道这个请求的作用域是哪个键,这就是为什么你会得到一个异常。

      另一方面,您可以安全地在processElementKeyedBroadcastProcessFunction 中使用键控状态,因为这些元素将分配键并且在键控状态的情况下范围是已知的。

      如果您需要对非广播状态的广播元素使用状态,则需要按照文档中的说明将其实现为操作员状态。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2021-12-29
        • 2020-07-27
        • 1970-01-01
        • 2022-11-30
        • 1970-01-01
        • 1970-01-01
        • 2017-07-11
        • 2020-08-17
        相关资源
        最近更新 更多